From 53598c8e1498e7288d957dda5b690b21d8f2d953 Mon Sep 17 00:00:00 2001 From: Manus AI Date: Sat, 12 Sep 2026 12:34:50 +0000 Subject: [PATCH] feat: add runtime audit adapters --- cmd/agentctl/main.go | 112 ++++++++++++- docs/runtime-audit-mvp.md | 103 ++++++++++++ internal/runtime/adapters.go | 254 ++++++++++++++++++++++++++++++ internal/runtime/adapters_test.go | 42 +++++ internal/runtime/aggregate.go | 78 +++++++++ internal/runtime/audit.go | 134 ++++++++++++++++ internal/runtime/audit_test.go | 37 +++++ internal/runtime/ingest.go | 72 +++++++++ internal/runtime/model.go | 43 +++++ internal/runtime/runtime_test.go | 56 +++++++ 10 files changed, 929 insertions(+), 2 deletions(-) create mode 100644 docs/runtime-audit-mvp.md create mode 100644 internal/runtime/adapters.go create mode 100644 internal/runtime/adapters_test.go create mode 100644 internal/runtime/aggregate.go create mode 100644 internal/runtime/audit.go create mode 100644 internal/runtime/audit_test.go create mode 100644 internal/runtime/ingest.go create mode 100644 internal/runtime/model.go create mode 100644 internal/runtime/runtime_test.go diff --git a/cmd/agentctl/main.go b/cmd/agentctl/main.go index 2fc1e8b..fc46a4c 100644 --- a/cmd/agentctl/main.go +++ b/cmd/agentctl/main.go @@ -12,6 +12,7 @@ import ( "time" "github.com/HlinorAI/agent-control-plane/internal/config" + "github.com/HlinorAI/agent-control-plane/internal/runtime" "github.com/HlinorAI/agent-control-plane/internal/scan" ) @@ -23,6 +24,7 @@ Usage: agentctl version agentctl init agentctl scan [flags] + agentctl runtime-audit [flags] Scan flags: --baseline file suppress findings already present in a JSON report @@ -32,7 +34,14 @@ Scan flags: --fail-on severity fail when findings meet severity: none, low, medium, high, critical --format format output format: text, json or sarif --output file write the report to a file instead of stdout - --suppressions file suppress active findings with reason and expiry from a JSON file + --suppressions file suppress active findings with reason and expiry from a JSON file + +Runtime audit flags: + --source format input source: jsonl, otel-json or api-gateway + --fail-on severity return non-zero at this severity or higher + --format format output format: text or json + --inventory file static agentctl JSON report to compare with runtime events + --output file write the audit report to a file instead of stdout The scanner is read-only and metadata-only. It does not execute scanned content. ` @@ -65,8 +74,11 @@ func run(args []string, stdout, stderr io.Writer) error { if args[0] == "init" { return runInit(args[1:], stdout, stderr) } + if args[0] == "runtime-audit" { + return runRuntimeAudit(args[1:], stdout, stderr) + } if args[0] != "scan" { - return fmt.Errorf("unknown command %q; available commands are init and scan", args[0]) + return fmt.Errorf("unknown command %q; available commands are init, scan and runtime-audit", args[0]) } if len(args) == 2 && (args[1] == "--help" || args[1] == "-h") { _, err := io.WriteString(stdout, usage) @@ -381,6 +393,102 @@ func runInit(args []string, stdout, stderr io.Writer) error { return err } +func runRuntimeAudit(args []string, stdout, stderr io.Writer) error { + if len(args) < 1 { + return errors.New("runtime-audit requires an events JSONL path") + } + fs := flag.NewFlagSet("runtime-audit", flag.ContinueOnError) + fs.SetOutput(stderr) + source := fs.String("source", string(runtime.SourceJSONL), "input source: jsonl, otel-json or api-gateway") + inventoryPath := fs.String("inventory", "", "static agentctl JSON report") + format := fs.String("format", "text", "output format: text or json") + failOn := fs.String("fail-on", "none", "return a non-zero exit code at this severity or higher") + output := fs.String("output", "", "write the audit report to a file instead of stdout") + if err := fs.Parse(args[1:]); err != nil { + return err + } + if *inventoryPath == "" { + return errors.New("runtime-audit requires --inventory") + } + if *source != string(runtime.SourceJSONL) && *source != string(runtime.SourceOTelJSON) && *source != string(runtime.SourceAPIGateway) { + return fmt.Errorf("unsupported runtime source %q", *source) + } + if *format != "text" && *format != "json" { + return fmt.Errorf("unsupported runtime audit format %q", *format) + } + if !validSeverity(*failOn) { + return fmt.Errorf("unsupported fail-on severity %q", *failOn) + } + eventsFile, err := os.Open(args[0]) + if err != nil { + return fmt.Errorf("open runtime events: %w", err) + } + defer eventsFile.Close() + events, skipped, err := runtime.ReadSource(eventsFile, runtime.Source(*source), runtime.Options{}) + if err != nil { + return err + } + inventoryFile, err := os.Open(*inventoryPath) + if err != nil { + return fmt.Errorf("open inventory: %w", err) + } + defer inventoryFile.Close() + var inventory scan.Report + decoder := json.NewDecoder(inventoryFile) + if err := decoder.Decode(&inventory); err != nil { + return fmt.Errorf("decode inventory: %w", err) + } + audit := runtime.Audit(runtime.Aggregate(events, skipped), inventory) + var payload []byte + if *format == "json" { + payload, err = json.MarshalIndent(audit, "", " ") + if err == nil { + payload = append(payload, '\n') + } + } else { + payload = []byte(runtimeAuditText(audit)) + } + if err != nil { + return err + } + if *output != "" { + if err := writeOutputFile(*output, payload); err != nil { + return err + } + } else if _, err := stdout.Write(payload); err != nil { + return err + } + if runtimeFindingsMeetThreshold(audit.Findings, *failOn) { + return fmt.Errorf("runtime audit found findings at or above %s severity", *failOn) + } + return nil +} + +func runtimeAuditText(report runtime.AuditReport) string { + var b strings.Builder + fmt.Fprintf(&b, "Agent Control Plane runtime audit\nSchema: %s\nEvents: %d\nMatched agents: %d\nUnmatched agents: %d\nFindings: %d\n", report.SchemaVersion, report.Runtime.EventsRead, report.MatchedAgents, report.Unmatched, len(report.Findings)) + for _, finding := range report.Findings { + fmt.Fprintf(&b, "- [%s] %s: %s\n", finding.Severity, finding.RuleID, finding.Message) + for _, evidence := range finding.Evidence { + fmt.Fprintf(&b, " evidence: %s\n", evidence) + } + } + return b.String() +} + +func runtimeFindingsMeetThreshold(findings []runtime.Finding, threshold string) bool { + minimum := severityRank(threshold) + if minimum == 0 { + return false + } + for _, finding := range findings { + if severityRank(finding.Severity) >= minimum { + return true + } + } + return false +} + func resolveConfigPath(root, requested string) (string, error) { absRoot, err := filepath.Abs(root) if err != nil { diff --git a/docs/runtime-audit-mvp.md b/docs/runtime-audit-mvp.md new file mode 100644 index 0000000..ac20e81 --- /dev/null +++ b/docs/runtime-audit-mvp.md @@ -0,0 +1,103 @@ +# Runtime-аудит Agent Control Plane: MVP + +## Решение + +Agent Control Plane следует расширять отдельным **runtime evidence layer**, а не превращать статический сканер в прокси или систему исполнения. Статический сканер отвечает на вопрос «что заявлено в коде», а runtime-аудит — на вопрос «что фактически происходило». Эти результаты должны связываться по стабильному идентификатору агента и объединяться только на этапе отчётности. + +MVP должен принимать нормализованные события в формате JSON Lines, агрегировать фактические вызовы и сопоставлять их с результатом статического сканирования. На первом этапе система не должна исполнять инструменты, перехватывать секреты, менять политики или требовать конкретного observability-провайдера. + +## Цели и ограничения + +| Область | Входит в MVP | Не входит в MVP | +|---|---|---| +| Источники | Нормализованный JSONL, API Gateway JSON/JSONL, OpenTelemetry JSON export | Прямое подключение ко всем вендорам | +| Формат | JSONL с нормализованным событием | Произвольный парсинг логов каждого продукта | +| Аналитика | Число вызовов, успешность, уникальные цели, first/last seen | Поведенческая ML-детекция | +| Сопоставление | `agent_id`, затем устойчивое имя с явным предупреждением | Неоднозначное автоматическое связывание без evidence | +| Безопасность | Metadata-only, лимиты размера, отказ от payload/arguments | Сбор prompt, tool arguments и raw secrets | +| Выход | JSON-отчёт и runtime findings | Enforcement и runtime proxy | + +## Нормализованное событие + +Каждая строка входного потока представляет одно событие. Payload запроса и ответа намеренно отсутствуют. Идентификаторы и имена ограничиваются метаданными, необходимыми для аудита. + +```json +{ + "timestamp": "2026-09-12T12:00:00Z", + "request_id": "req-123", + "agent_id": "agent_customer_support", + "agent_name": "customer-support", + "environment": "production", + "operation": "tool_call", + "target": "crm.search_customers", + "provider": "internal-crm", + "action": "read", + "success": true +} +``` + +Обязательные поля: `timestamp`, `agent_id` или `agent_name`, `operation`, `target`, `success`. Значения `operation` и `action` являются свободными строками на этапе MVP, чтобы не блокировать интеграции. Нормализатор должен отклонять пустые идентификаторы и некорректные timestamps. + +## Runtime-агрегат + +Для каждого агента агрегируются следующие показатели: + +| Показатель | Назначение | +|---|---| +| `event_count` | Общий объём наблюдений | +| `successful_events` и `failed_events` | Надёжность и ошибки интеграций | +| `unique_targets` | Фактический scope инструментов и API | +| `operations` | Разбивка по типам операций | +| `providers` | Фактические model/API providers | +| `first_seen` и `last_seen` | Окно наблюдения | +| `undeclared_targets` | Цели, отсутствующие в статическом inventory | +| `undeclared_providers` | Провайдеры, отсутствующие в декларации | + +Агрегаты сортируются детерминированно. Это необходимо для стабильных diff-отчётов и baseline-механизма, уже используемого статическим сканером. + +## Runtime findings + +Первый набор правил должен быть небольшим и проверяемым: + +| Правило | Условие | Начальная severity | +|---|---|---| +| `ACP-R001` | Фактическая цель отсутствует в статическом inventory агента | High | +| `ACP-R002` | Фактический provider отличается от заявленного | High | +| `ACP-R003` | Production-агент обращается к цели с write-действием, хотя заявлен как read-only | Critical | +| `ACP-R004` | Runtime-события невозможно надёжно сопоставить с агентом | Medium | +| `ACP-R005` | Наблюдаемые события старше заданного окна свежести | Note | + +Правила должны создавать evidence только из безопасных полей: `request_id`, timestamp, operation и target. Никогда не следует помещать в finding prompt, аргументы инструмента, заголовки авторизации или тело ответа. + +## План реализации + +1. **Слой ingest.** Добавить потоковый JSONL reader с лимитом строки и лимитом общего числа событий. +2. **Агрегация.** Добавить детерминированный runtime report с first/last seen и уникальными целями. +3. **Сопоставление.** Сопоставить runtime agent с `scan.Report.Agents`; при отсутствии `agent_id` разрешать имя только при единственном совпадении. +4. **Findings.** Добавить первые runtime rules без enforcement. +5. **CLI.** Добавить отдельную команду `agentctl runtime-audit --inventory ` после стабилизации библиотечного API. +6. **Интеграции.** Реализовать адаптеры для OpenTelemetry и API gateway после утверждения нормализованной схемы. + +## Критерии готовности MVP + +MVP считается готовым, когда он обрабатывает поток минимум из 100 000 metadata-only событий с ограниченным потреблением памяти, не выводит запрещённые payload-поля, выдаёт одинаковый JSON для одинакового входа, корректно обрабатывает повреждённые строки с диагностикой и создаёт отдельные findings для undeclared target/provider. + +## Принцип безопасности + +Runtime-аудит должен быть **наблюдателем, а не исполнителем**. Он не запускает команды из событий, не обращается к URL из полей `target`, не читает содержимое tool arguments и не принимает решения о выдаче доступа. Enforcement — отдельный будущий компонент с самостоятельной моделью угроз и процессом согласования. + +## Реализовано в текущем этапе + +В репозитории реализован библиотечный слой `internal/runtime`, который выполняет безопасное чтение нормализованных JSONL-событий, детерминированную агрегацию и сопоставление со статическим `scan.Report`. Добавлены правила `ACP-R001` — `ACP-R004`, включая обнаружение undeclared target, фактического provider, production write activity и неоднозначного сопоставления агента. + +Добавлена команда `agentctl runtime-audit --inventory `. Флаг `--source` выбирает `jsonl`, `otel-json` или `api-gateway`. Команда поддерживает text/json output, атомарную запись через существующий механизм `--output` и CI-порог `--fail-on`. + +Адаптеры принимают только metadata-поля. OpenTelemetry spans преобразуются по атрибутам агента, инструмента, окружения, provider и HTTP-статуса. API Gateway записи поддерживают snake_case и camelCase идентификаторы, JSONL и JSON-массив. Payload запроса и ответа не читается и не переносится в нормализованное событие. + +Следующий этап — добавить адаптеры OpenTelemetry и API gateway, не меняя нормализованный контракт событий. + +## References + +[1]: https://opentelemetry.io/docs/specs/otel/ "OpenTelemetry Specification" + +[2]: https://www.aicpa-cima.com/resources/download/2017-trust-services-criteria-tsp-section-100 "AICPA Trust Services Criteria" diff --git a/internal/runtime/adapters.go b/internal/runtime/adapters.go new file mode 100644 index 0000000..107c146 --- /dev/null +++ b/internal/runtime/adapters.go @@ -0,0 +1,254 @@ +package runtime + +import ( + "bufio" + "bytes" + "encoding/json" + "fmt" + "io" + "strconv" + "strings" + "time" +) + +type Source string + +const maxAdapterBytes = 64 << 20 + +const ( + SourceJSONL Source = "jsonl" + SourceOTelJSON Source = "otel-json" + SourceAPIGateway Source = "api-gateway" +) + +func ReadSource(r io.Reader, source Source, options Options) ([]Event, int, error) { + switch source { + case "", SourceJSONL: + return ReadEvents(r, options) + case SourceOTelJSON: + return readOTelJSON(r, options) + case SourceAPIGateway: + return readAPIGateway(r, options) + default: + return nil, 0, fmt.Errorf("unsupported runtime source %q", source) + } +} + +type otelDocument struct { + ResourceSpans []otelResourceSpans `json:"resourceSpans"` +} + +type otelResourceSpans struct { + Resource otelResource `json:"resource"` + ScopeSpans []otelScopeSpans `json:"scopeSpans"` +} + +type otelScopeSpans struct { + Spans []otelSpan `json:"spans"` +} + +type otelResource struct { + Attributes []otelAttribute `json:"attributes"` +} + +type otelSpan struct { + TraceID string `json:"traceId"` + Name string `json:"name"` + StartTimeUnixNano json.Number `json:"startTimeUnixNano"` + StartTime string `json:"startTime"` + Status otelStatus `json:"status"` + Attributes []otelAttribute `json:"attributes"` +} + +type otelStatus struct { + Code string `json:"code"` +} + +type otelAttribute struct { + Key string `json:"key"` + Value otelAttrValue `json:"value"` +} + +type otelAttrValue struct { + StringValue string `json:"stringValue"` + IntValue json.Number `json:"intValue"` + BoolValue *bool `json:"boolValue"` +} + +func readOTelJSON(r io.Reader, options Options) ([]Event, int, error) { + options = options.withDefaults() + var document otelDocument + decoder := json.NewDecoder(io.LimitReader(r, maxAdapterBytes)) + if err := decoder.Decode(&document); err != nil { + return nil, 0, fmt.Errorf("decode OpenTelemetry JSON: %w", err) + } + events := make([]Event, 0) + for _, resourceSpans := range document.ResourceSpans { + resourceAttrs := attributes(resourceSpans.Resource.Attributes) + for _, scopeSpans := range resourceSpans.ScopeSpans { + for _, span := range scopeSpans.Spans { + if len(events) >= options.MaxEvents { + return nil, 0, fmt.Errorf("runtime event limit exceeded: %d", options.MaxEvents) + } + event, err := eventFromOTelSpan(span, resourceAttrs) + if err != nil { + return nil, 0, err + } + events = append(events, event) + } + } + } + return events, 0, nil +} + +func eventFromOTelSpan(span otelSpan, resourceAttrs map[string]string) (Event, error) { + attrs := attributes(span.Attributes) + for key, value := range resourceAttrs { + if _, exists := attrs[key]; !exists { + attrs[key] = value + } + } + timestamp, err := parseOTelTime(span.StartTime, span.StartTimeUnixNano) + if err != nil { + return Event{}, fmt.Errorf("decode OpenTelemetry span %q timestamp: %w", span.Name, err) + } + target := firstNonEmpty(attrs["agent.target"], attrs["tool.name"], attrs["rpc.method"], attrs["http.route"], attrs["server.address"], span.Name) + operation := firstNonEmpty(attrs["agent.operation"], attrs["operation"], "trace:"+span.Name) + success := !strings.EqualFold(span.Status.Code, "ERROR") && !strings.EqualFold(attrs["otel.status_code"], "ERROR") + if statusCode, parseErr := strconv.Atoi(attrs["http.status_code"]); parseErr == nil && statusCode >= 400 { + success = false + } + event := Event{ + Timestamp: timestamp, + RequestID: firstNonEmpty(attrs["request.id"], attrs["http.request_id"], span.TraceID), + AgentID: firstNonEmpty(attrs["agent.id"], attrs["gen_ai.agent.id"]), + AgentName: firstNonEmpty(attrs["agent.name"], attrs["gen_ai.agent.name"]), + Environment: firstNonEmpty(attrs["deployment.environment"], attrs["deployment.environment.name"]), + Operation: operation, + Target: target, + Provider: firstNonEmpty(attrs["gen_ai.system"], attrs["gen_ai.provider.name"], attrs["server.address"]), + Action: firstNonEmpty(attrs["agent.action"], attrs["tool.action"]), + Success: success, + } + if err := validateEvent(event); err != nil { + return Event{}, fmt.Errorf("validate OpenTelemetry span %q: %w", span.Name, err) + } + return event, nil +} + +func parseOTelTime(value string, unixNano json.Number) (time.Time, error) { + if value != "" { + parsed, err := time.Parse(time.RFC3339Nano, value) + if err != nil { + return time.Time{}, err + } + return parsed, nil + } + nanos, err := strconv.ParseInt(string(unixNano), 10, 64) + if err != nil || nanos <= 0 { + return time.Time{}, fmt.Errorf("invalid startTimeUnixNano") + } + return time.Unix(0, nanos).UTC(), nil +} + +func attributes(values []otelAttribute) map[string]string { + result := make(map[string]string, len(values)) + for _, attribute := range values { + value := attribute.Value.StringValue + if value == "" { + value = string(attribute.Value.IntValue) + } + if value == "" && attribute.Value.BoolValue != nil { + value = strconv.FormatBool(*attribute.Value.BoolValue) + } + if attribute.Key != "" && value != "" { + result[strings.ToLower(attribute.Key)] = value + } + } + return result +} + +type gatewayRecord struct { + Timestamp string `json:"timestamp"` + RequestID string `json:"request_id"` + RequestID2 string `json:"requestId"` + AgentID string `json:"agent_id"` + AgentID2 string `json:"agentId"` + AgentName string `json:"agent_name"` + Environment string `json:"environment"` + Operation string `json:"operation"` + Target string `json:"target"` + Provider string `json:"provider"` + Action string `json:"action"` + Success *bool `json:"success"` + Status string `json:"status"` + StatusCode int `json:"status_code"` +} + +func readAPIGateway(r io.Reader, options Options) ([]Event, int, error) { + options = options.withDefaults() + data, err := io.ReadAll(io.LimitReader(r, maxAdapterBytes)) + if err != nil { + return nil, 0, fmt.Errorf("read API Gateway events: %w", err) + } + data = bytes.TrimSpace(data) + if len(data) == 0 { + return nil, 0, nil + } + if data[0] == '[' { + var records []gatewayRecord + if err := json.Unmarshal(data, &records); err != nil { + return nil, 0, fmt.Errorf("decode API Gateway JSON: %w", err) + } + return gatewayEvents(records, options) + } + scanner := bufio.NewScanner(bytes.NewReader(data)) + scanner.Buffer(make([]byte, 64*1024), options.MaxLineBytes) + records := make([]gatewayRecord, 0) + for scanner.Scan() { + if strings.TrimSpace(scanner.Text()) == "" { + continue + } + var record gatewayRecord + if err := json.Unmarshal(scanner.Bytes(), &record); err != nil { + return nil, 0, fmt.Errorf("decode API Gateway JSONL event %d: %w", len(records)+1, err) + } + records = append(records, record) + } + if err := scanner.Err(); err != nil { + return nil, 0, fmt.Errorf("read API Gateway JSONL: %w", err) + } + return gatewayEvents(records, options) +} + +func gatewayEvents(records []gatewayRecord, options Options) ([]Event, int, error) { + if len(records) > options.MaxEvents { + return nil, 0, fmt.Errorf("runtime event limit exceeded: %d", options.MaxEvents) + } + events := make([]Event, 0, len(records)) + for index, record := range records { + timestamp, err := time.Parse(time.RFC3339Nano, record.Timestamp) + if err != nil { + return nil, 0, fmt.Errorf("decode API Gateway event %d timestamp: %w", index+1, err) + } + success := record.Success == nil || *record.Success + if record.StatusCode >= 400 || strings.EqualFold(record.Status, "error") || strings.EqualFold(record.Status, "failed") { + success = false + } + event := Event{Timestamp: timestamp, RequestID: firstNonEmpty(record.RequestID, record.RequestID2), AgentID: firstNonEmpty(record.AgentID, record.AgentID2), AgentName: record.AgentName, Environment: record.Environment, Operation: record.Operation, Target: record.Target, Provider: record.Provider, Action: record.Action, Success: success} + if err := validateEvent(event); err != nil { + return nil, 0, fmt.Errorf("validate API Gateway event %d: %w", index+1, err) + } + events = append(events, event) + } + return events, 0, nil +} + +func firstNonEmpty(values ...string) string { + for _, value := range values { + if strings.TrimSpace(value) != "" { + return value + } + } + return "" +} diff --git a/internal/runtime/adapters_test.go b/internal/runtime/adapters_test.go new file mode 100644 index 0000000..449b10d --- /dev/null +++ b/internal/runtime/adapters_test.go @@ -0,0 +1,42 @@ +package runtime + +import ( + "strings" + "testing" +) + +func TestReadOTelJSON(t *testing.T) { + input := `{"resourceSpans":[{"resource":{"attributes":[{"key":"deployment.environment","value":{"stringValue":"production"}}]},"scopeSpans":[{"spans":[{"traceId":"trace-1","name":"crm.search","startTime":"2026-01-02T03:04:05Z","status":{"code":"OK"},"attributes":[{"key":"agent.id","value":{"stringValue":"a1"}},{"key":"agent.name","value":{"stringValue":"support"}},{"key":"tool.name","value":{"stringValue":"crm.search"}},{"key":"gen_ai.system","value":{"stringValue":"declared-provider"}}]}]}]}]}` + events, skipped, err := ReadSource(strings.NewReader(input), SourceOTelJSON, Options{}) + if err != nil { + t.Fatalf("ReadSource() error = %v", err) + } + if skipped != 0 || len(events) != 1 { + t.Fatalf("unexpected events: skipped=%d events=%+v", skipped, events) + } + event := events[0] + if event.AgentID != "a1" || event.Target != "crm.search" || event.Environment != "production" || !event.Success { + t.Fatalf("unexpected normalized OTel event: %+v", event) + } +} + +func TestReadAPIGatewayJSONL(t *testing.T) { + input := `{"timestamp":"2026-01-02T03:04:05Z","requestId":"req-1","agentId":"a1","agent_name":"support","environment":"production","operation":"http_request","target":"/customers","provider":"crm","action":"read","status_code":200} +{"timestamp":"2026-01-02T03:05:05Z","agent_id":"a1","operation":"http_request","target":"/customers","status":"error","success":true}` + events, skipped, err := ReadSource(strings.NewReader(input), SourceAPIGateway, Options{}) + if err != nil { + t.Fatalf("ReadSource() error = %v", err) + } + if skipped != 0 || len(events) != 2 { + t.Fatalf("unexpected events: skipped=%d events=%+v", skipped, events) + } + if events[0].RequestID != "req-1" || events[0].AgentID != "a1" || events[1].Success { + t.Fatalf("unexpected normalized gateway events: %+v", events) + } +} + +func TestReadSourceRejectsUnknownSource(t *testing.T) { + if _, _, err := ReadSource(strings.NewReader("{}"), Source("unknown"), Options{}); err == nil { + t.Fatal("ReadSource() accepted unknown source") + } +} diff --git a/internal/runtime/aggregate.go b/internal/runtime/aggregate.go new file mode 100644 index 0000000..aa7d690 --- /dev/null +++ b/internal/runtime/aggregate.go @@ -0,0 +1,78 @@ +package runtime + +import ( + "sort" + "strings" +) + +func Aggregate(events []Event, skipped int) Report { + byKey := make(map[string]*AgentSummary) + for _, event := range events { + key := event.AgentID + if key == "" { + key = "name:" + event.AgentName + } + summary := byKey[key] + if summary == nil { + summary = &AgentSummary{ + AgentID: event.AgentID, + AgentName: event.AgentName, + FirstSeen: event.Timestamp, + LastSeen: event.Timestamp, + Operations: make(map[string]int), + } + byKey[key] = summary + } + summary.EventCount++ + if event.Success { + summary.SuccessfulEvents++ + } else { + summary.FailedEvents++ + } + if strings.EqualFold(event.Action, "write") || strings.EqualFold(event.Action, "delete") || strings.EqualFold(event.Action, "mutate") { + summary.WriteEvents++ + } + if event.Timestamp.Before(summary.FirstSeen) { + summary.FirstSeen = event.Timestamp + } + if event.Timestamp.After(summary.LastSeen) { + summary.LastSeen = event.Timestamp + } + summary.Targets = appendUnique(summary.Targets, event.Target) + summary.Environments = appendUnique(summary.Environments, event.Environment) + summary.Providers = appendUnique(summary.Providers, event.Provider) + summary.Operations[event.Operation]++ + } + + agents := make([]AgentSummary, 0, len(byKey)) + for _, summary := range byKey { + sort.Strings(summary.Targets) + sort.Strings(summary.Environments) + sort.Strings(summary.Providers) + agents = append(agents, *summary) + } + sort.Slice(agents, func(i, j int) bool { + left := agents[i].AgentID + if left == "" { + left = "name:" + agents[i].AgentName + } + right := agents[j].AgentID + if right == "" { + right = "name:" + agents[j].AgentName + } + return left < right + }) + return Report{SchemaVersion: "runtime.v1", EventsRead: len(events), EventsSkipped: skipped, Agents: agents} +} + +func appendUnique(values []string, value string) []string { + if strings.TrimSpace(value) == "" { + return values + } + for _, current := range values { + if current == value { + return values + } + } + return append(values, value) +} diff --git a/internal/runtime/audit.go b/internal/runtime/audit.go new file mode 100644 index 0000000..9923773 --- /dev/null +++ b/internal/runtime/audit.go @@ -0,0 +1,134 @@ +package runtime + +import ( + "crypto/sha256" + "fmt" + "sort" + "strings" + + "github.com/HlinorAI/agent-control-plane/internal/scan" +) + +type Finding struct { + ID string `json:"id"` + RuleID string `json:"rule_id"` + Severity string `json:"severity"` + Message string `json:"message"` + AgentID string `json:"agent_id,omitempty"` + Confidence float64 `json:"confidence"` + Evidence []string `json:"evidence,omitempty"` + RemediationHint string `json:"remediation_hint"` +} + +type AuditReport struct { + SchemaVersion string `json:"schema_version"` + Runtime Report `json:"runtime"` + MatchedAgents int `json:"matched_agents"` + Unmatched int `json:"unmatched_agents"` + Findings []Finding `json:"findings"` +} + +func Audit(runtimeReport Report, inventory scan.Report) AuditReport { + result := AuditReport{SchemaVersion: "runtime-audit.v1", Runtime: runtimeReport} + byID := make(map[string]scan.Agent, len(inventory.Agents)) + byName := make(map[string][]scan.Agent) + for _, agent := range inventory.Agents { + byID[agent.ID] = agent + name := strings.ToLower(strings.TrimSpace(agent.Name)) + if name != "" { + byName[name] = append(byName[name], agent) + } + } + providers := make(map[string]bool) + for _, model := range inventory.Models { + providers[strings.ToLower(strings.TrimSpace(model.Provider))] = true + } + + for _, summary := range runtimeReport.Agents { + agent, ok := matchAgent(summary, byID, byName) + if !ok { + result.Unmatched++ + result.Findings = append(result.Findings, Finding{ + RuleID: "ACP-R004", Severity: "Medium", Message: fmt.Sprintf("runtime activity for %q cannot be matched uniquely to static inventory", runtimeAgentLabel(summary)), Confidence: 0.92, + Evidence: []string{summary.FirstSeen.Format("2006-01-02T15:04:05Z07:00"), summary.LastSeen.Format("2006-01-02T15:04:05Z07:00")}, RemediationHint: "Emit the stable static agent ID in runtime telemetry.", + }) + continue + } + result.MatchedAgents++ + declaredTargets := make(map[string]bool) + for _, target := range agent.Tools { + declaredTargets[strings.ToLower(strings.TrimSpace(target))] = true + } + for _, target := range summary.Targets { + if !declaredTargets[strings.ToLower(strings.TrimSpace(target))] { + result.Findings = append(result.Findings, Finding{ + RuleID: "ACP-R001", Severity: "High", Message: fmt.Sprintf("agent %q used undeclared runtime target %q", agent.Name, target), AgentID: agent.ID, Confidence: 0.90, + Evidence: []string{target, summary.LastSeen.Format("2006-01-02T15:04:05Z07:00")}, RemediationHint: "Declare the target in the agent inventory or investigate the unexpected integration.", + }) + } + } + for _, provider := range summary.Providers { + if provider != "" && !providers[strings.ToLower(strings.TrimSpace(provider))] { + result.Findings = append(result.Findings, Finding{ + RuleID: "ACP-R002", Severity: "High", Message: fmt.Sprintf("agent %q used provider %q not present in static model inventory", agent.Name, provider), AgentID: agent.ID, Confidence: 0.86, + Evidence: []string{provider, summary.LastSeen.Format("2006-01-02T15:04:05Z07:00")}, RemediationHint: "Verify the provider and update workspace policy or runtime configuration.", + }) + } + } + if summary.WriteEvents > 0 && hasProduction(summary.Environments) { + undeclaredWrite := false + for _, target := range summary.Targets { + if !declaredTargets[strings.ToLower(strings.TrimSpace(target))] { + undeclaredWrite = true + } + } + if undeclaredWrite { + result.Findings = append(result.Findings, Finding{ + RuleID: "ACP-R003", Severity: "Critical", Message: fmt.Sprintf("production agent %q performed write activity against an undeclared target", agent.Name), AgentID: agent.ID, Confidence: 0.94, + Evidence: []string{summary.LastSeen.Format("2006-01-02T15:04:05Z07:00")}, RemediationHint: "Block or review the write path and declare an approved least-privilege scope.", + }) + } + } + } + sort.Slice(result.Findings, func(i, j int) bool { + left := result.Findings[i].RuleID + ":" + result.Findings[i].AgentID + ":" + result.Findings[i].Message + right := result.Findings[j].RuleID + ":" + result.Findings[j].AgentID + ":" + result.Findings[j].Message + return left < right + }) + for i := range result.Findings { + result.Findings[i].ID = findingID(result.Findings[i]) + } + return result +} + +func matchAgent(summary AgentSummary, byID map[string]scan.Agent, byName map[string][]scan.Agent) (scan.Agent, bool) { + if agent, ok := byID[summary.AgentID]; ok && summary.AgentID != "" { + return agent, true + } + matches := byName[strings.ToLower(strings.TrimSpace(summary.AgentName))] + if len(matches) == 1 { + return matches[0], true + } + return scan.Agent{}, false +} + +func runtimeAgentLabel(summary AgentSummary) string { + if summary.AgentName != "" { + return summary.AgentName + } + return summary.AgentID +} + +func hasProduction(environments []string) bool { + for _, environment := range environments { + if strings.EqualFold(environment, "production") || strings.EqualFold(environment, "prod") { + return true + } + } + return false +} + +func findingID(finding Finding) string { + sum := sha256.Sum256([]byte(finding.RuleID + ":" + finding.AgentID + ":" + finding.Message)) + return fmt.Sprintf("runtime_%x", sum[:8]) +} diff --git a/internal/runtime/audit_test.go b/internal/runtime/audit_test.go new file mode 100644 index 0000000..21822cb --- /dev/null +++ b/internal/runtime/audit_test.go @@ -0,0 +1,37 @@ +package runtime + +import ( + "testing" + "time" + + "github.com/HlinorAI/agent-control-plane/internal/scan" +) + +func TestAuditFindsUndeclaredProductionWrite(t *testing.T) { + runtimeReport := Aggregate([]Event{{ + Timestamp: time.Date(2026, 1, 2, 3, 4, 5, 0, time.UTC), AgentID: "a1", AgentName: "support", Environment: "production", Operation: "tool_call", Target: "crm.write", Provider: "unlisted", Action: "write", Success: true, + }}, 0) + result := Audit(runtimeReport, scan.Report{Agents: []scan.Agent{{ID: "a1", Name: "support", Tools: []string{"crm.search"}}}}) + if result.MatchedAgents != 1 || len(result.Findings) != 3 { + t.Fatalf("unexpected audit result: %+v", result) + } + seen := map[string]bool{} + for _, finding := range result.Findings { + seen[finding.RuleID] = true + } + for _, rule := range []string{"ACP-R001", "ACP-R002", "ACP-R003"} { + if !seen[rule] { + t.Fatalf("missing finding %s: %+v", rule, result.Findings) + } + } +} + +func TestAuditRejectsAmbiguousNameMatch(t *testing.T) { + runtimeReport := Aggregate([]Event{{ + Timestamp: time.Date(2026, 1, 2, 3, 4, 5, 0, time.UTC), AgentName: "support", Operation: "tool_call", Target: "crm.search", Success: true, + }}, 0) + result := Audit(runtimeReport, scan.Report{Agents: []scan.Agent{{ID: "a1", Name: "support"}, {ID: "a2", Name: "support"}}}) + if result.MatchedAgents != 0 || result.Unmatched != 1 || len(result.Findings) != 1 || result.Findings[0].RuleID != "ACP-R004" { + t.Fatalf("unexpected ambiguous match result: %+v", result) + } +} diff --git a/internal/runtime/ingest.go b/internal/runtime/ingest.go new file mode 100644 index 0000000..d60c16a --- /dev/null +++ b/internal/runtime/ingest.go @@ -0,0 +1,72 @@ +package runtime + +import ( + "bufio" + "encoding/json" + "fmt" + "io" + "strings" + "time" +) + +type Options struct { + MaxLineBytes int + MaxEvents int +} + +func (o Options) withDefaults() Options { + if o.MaxLineBytes <= 0 { + o.MaxLineBytes = DefaultMaxLineBytes + } + if o.MaxEvents <= 0 { + o.MaxEvents = DefaultMaxEvents + } + return o +} + +func ReadEvents(r io.Reader, options Options) ([]Event, int, error) { + options = options.withDefaults() + scanner := bufio.NewScanner(r) + scanner.Buffer(make([]byte, 64*1024), options.MaxLineBytes) + events := make([]Event, 0) + skipped := 0 + for scanner.Scan() { + if len(scanner.Bytes()) == 0 || strings.TrimSpace(scanner.Text()) == "" { + continue + } + if len(events) >= options.MaxEvents { + return nil, skipped, fmt.Errorf("runtime event limit exceeded: %d", options.MaxEvents) + } + var event Event + if err := json.Unmarshal(scanner.Bytes(), &event); err != nil { + return nil, skipped, fmt.Errorf("decode runtime event %d: %w", len(events)+skipped+1, err) + } + if err := validateEvent(event); err != nil { + return nil, skipped, fmt.Errorf("validate runtime event %d: %w", len(events)+skipped+1, err) + } + events = append(events, event) + } + if err := scanner.Err(); err != nil { + return nil, skipped, fmt.Errorf("read runtime events: %w", err) + } + return events, skipped, nil +} + +func validateEvent(event Event) error { + if event.Timestamp.IsZero() { + return fmt.Errorf("timestamp is required") + } + if event.Timestamp.After(time.Now().Add(5 * time.Minute)) { + return fmt.Errorf("timestamp is in the future") + } + if strings.TrimSpace(event.AgentID) == "" && strings.TrimSpace(event.AgentName) == "" { + return fmt.Errorf("agent_id or agent_name is required") + } + if strings.TrimSpace(event.Operation) == "" { + return fmt.Errorf("operation is required") + } + if strings.TrimSpace(event.Target) == "" { + return fmt.Errorf("target is required") + } + return nil +} diff --git a/internal/runtime/model.go b/internal/runtime/model.go new file mode 100644 index 0000000..af15eca --- /dev/null +++ b/internal/runtime/model.go @@ -0,0 +1,43 @@ +package runtime + +import "time" + +const ( + DefaultMaxLineBytes = 1 << 20 + DefaultMaxEvents = 1_000_000 +) + +type Event struct { + Timestamp time.Time `json:"timestamp"` + RequestID string `json:"request_id,omitempty"` + AgentID string `json:"agent_id,omitempty"` + AgentName string `json:"agent_name,omitempty"` + Environment string `json:"environment,omitempty"` + Operation string `json:"operation"` + Target string `json:"target"` + Provider string `json:"provider,omitempty"` + Action string `json:"action,omitempty"` + Success bool `json:"success"` +} + +type AgentSummary struct { + AgentID string `json:"agent_id,omitempty"` + AgentName string `json:"agent_name,omitempty"` + EventCount int `json:"event_count"` + SuccessfulEvents int `json:"successful_events"` + FailedEvents int `json:"failed_events"` + WriteEvents int `json:"write_events"` + FirstSeen time.Time `json:"first_seen"` + LastSeen time.Time `json:"last_seen"` + Targets []string `json:"targets,omitempty"` + Environments []string `json:"environments,omitempty"` + Operations map[string]int `json:"operations,omitempty"` + Providers []string `json:"providers,omitempty"` +} + +type Report struct { + SchemaVersion string `json:"schema_version"` + EventsRead int `json:"events_read"` + EventsSkipped int `json:"events_skipped"` + Agents []AgentSummary `json:"agents"` +} diff --git a/internal/runtime/runtime_test.go b/internal/runtime/runtime_test.go new file mode 100644 index 0000000..a33a437 --- /dev/null +++ b/internal/runtime/runtime_test.go @@ -0,0 +1,56 @@ +package runtime + +import ( + "strings" + "testing" + "time" +) + +func TestReadEventsAndAggregate(t *testing.T) { + input := `{"timestamp":"2026-01-02T03:04:05Z","request_id":"r1","agent_id":"a1","operation":"tool_call","target":"crm.search","provider":"crm","action":"read","success":true} +{"timestamp":"2026-01-02T03:05:05Z","request_id":"r2","agent_id":"a1","operation":"tool_call","target":"crm.write","provider":"crm","action":"write","success":false} +` + events, skipped, err := ReadEvents(strings.NewReader(input), Options{}) + if err != nil { + t.Fatalf("ReadEvents() error = %v", err) + } + report := Aggregate(events, skipped) + if report.SchemaVersion != "runtime.v1" || report.EventsRead != 2 || len(report.Agents) != 1 { + t.Fatalf("unexpected report: %+v", report) + } + summary := report.Agents[0] + if summary.EventCount != 2 || summary.SuccessfulEvents != 1 || summary.FailedEvents != 1 { + t.Fatalf("unexpected summary: %+v", summary) + } + if len(summary.Targets) != 2 || summary.Targets[0] != "crm.search" || summary.Targets[1] != "crm.write" { + t.Fatalf("targets are not sorted: %+v", summary.Targets) + } +} + +func TestReadEventsRejectsPayloadAndInvalidEvent(t *testing.T) { + input := `{"timestamp":"2026-01-02T03:04:05Z","operation":"tool_call","target":"crm.search","success":true}` + if _, _, err := ReadEvents(strings.NewReader(input), Options{}); err == nil { + t.Fatal("ReadEvents() accepted event without agent identity") + } +} + +func TestReadEventsLimit(t *testing.T) { + input := `{"timestamp":"2026-01-02T03:04:05Z","agent_name":"support","operation":"tool_call","target":"crm.search","success":true} +{"timestamp":"2026-01-02T03:04:06Z","agent_name":"support","operation":"tool_call","target":"crm.search","success":true} +` + _, _, err := ReadEvents(strings.NewReader(input), Options{MaxEvents: 1}) + if err == nil { + t.Fatal("ReadEvents() did not enforce event limit") + } +} + +func TestAggregateUsesEarliestAndLatestTimestamp(t *testing.T) { + events := []Event{ + {Timestamp: time.Date(2026, 1, 2, 3, 5, 0, 0, time.UTC), AgentName: "support", Operation: "tool_call", Target: "b", Success: true}, + {Timestamp: time.Date(2026, 1, 2, 3, 4, 0, 0, time.UTC), AgentName: "support", Operation: "model_call", Target: "a", Success: true}, + } + report := Aggregate(events, 0) + if !report.Agents[0].FirstSeen.Before(report.Agents[0].LastSeen) { + t.Fatalf("timestamps were not aggregated: %+v", report.Agents[0]) + } +}