agentleFS
Sign inSign up

khi-parser

GoogleCloudPlatform/khi/.agents/skills/khi_parser/SKILL.md

Guidelines, package patterns, and task implementations for adding new log type support or modifying existing log parsers in KHI.

Skill2.1k starsChanged 16 days ago
---
name: khi-parser
description: Guidelines, package patterns, and task implementations for adding new log type support or modifying existing log parsers in KHI.
---

# KHI Log Parser Support Guidelines

This guide outlines the patterns, package boundaries, implementation steps, and best practices for adding support for new log types or modifying existing log parsers in KHI.

---

## 1. Package Structure & Boundaries

When implementing a new log parser or modifying an existing one, you MUST separate the **contract** (IDs, public types, and configurations) from the **implementation** (the actual task logic). This guarantees that task IDs are fully initialized before implementation and prevents circular import dependencies.

The parser package must reside under `pkg/task/inspection/<provider>/<feature>/` (e.g., `pkg/task/inspection/googlecloud/k8snode/`) and adhere to the following structure:

```plaintext
pkg/task/inspection/<provider>/<feature>/
├── taskid.go          // Defines all TaskIDs and TaskReferences (package <feature>).
├── extractor.go       // (Optional) Defines field extraction functions and strongly-typed data structs.
├── timeline_type.go   // (Optional) Defines timeline types and verb types.
├── timeline_path.go   // (Optional) Helper functions to build hierarchical paths.
├── log_type.go        // (Optional) Defines log-specific types or constants.
└── impl/              // Package <feature>_impl
    ├── form.go            // (Optional) Implements form-related parameter tasks.
    ├── query.go           // (Optional) Implements log query/filter tasks.
    ├── ingester.go        // (Optional) Implements the LogIngester task.
    ├── grouper.go         // (Optional) Implements the LogGrouper task.
    ├── mapper.go          // (Optional) Implements LogToTimelineMapper tasks (or mapper_<target>.go).
    └── registration.go    // Implements task registration to the KHI registry.
```

### Key Package Boundaries

> [!IMPORTANT]
>
> - **Contract / Root Package (`pkg/task/inspection/<provider>/<feature>`)**: Uses `package <feature>` (no `_contract` suffix). MUST NOT import the `impl` package. External packages can freely import the root package to depend on parser task IDs, Extractor functions, or TimelineType constants.
> - **Implementation Package (`impl/`)**: Uses `package <feature>_impl`. Implements the actual tasks and imports the parent contract package. **External packages MUST NOT import the `impl` package.**
> - **File Naming**: Do NOT append redundant `_task.go` or `_tasks.go` suffixes to filenames in `impl/`. Name files strictly by their DAG pipeline role (`form.go`, `query.go`, `ingester.go`, `grouper.go`, `mapper.go` or `mapper_<target>.go`, `registration.go`).
> - **Registration**: Tasks inside the `impl` package are registered through `impl/registration.go`.

---

## 2. The Log Parsing Steps

A complete log parser in KHI generally consists of distinct DAG tasks:

```mermaid
flowchart TD
    FormTask[1. Form Task] -->|Provides Parameters| QueryTask[2. Log Query Task]
    QueryTask -->|Provides Raw Logs| IngesterTask[3. Log Ingest Task]
    QueryTask -->|Provides Raw Logs| GrouperTask[Log Grouper Task]
    IngesterTask -->|Provides Ingested Logs| MapperTask[4. Timeline Mapper Task]
    GrouperTask -->|Provides Grouped Map| MapperTask
```

### Step 1: Form Tasks (Form-related)

Exposes interactive input fields (e.g., text boxes, multi-select checkboxes) to let users configure parameters before running the inspection.

- **Utility:** `formtask.NewTextFormTaskBuilder` or `formtask.NewSetFormTaskBuilder`.

### Step 2: Log Query Tasks

Queries logs from the data source (e.g., Google Cloud Logging or local files) using parameters provided by the Form tasks.

- **Utility:** `gcpcommon.NewListLogEntriesTask` (for any logs on Cloud Logging) or `inspection_task.NewInspectionTask`.
- **Google Cloud API Calling:** When calling Google Cloud APIs directly or through fetchers, refer to [googlecloud-api](skill://googlecloud-api) for mandatory `CallOptionInjector` usage and client configuration.

### Step 3: Log Ingestion Tasks

Extracts information directly from the log's `NodeReader` using Extractor functions, populating basic log metadata on `LogChangeSet` (such as `Timestamp` from `l.Timestamp`, `Severity`, `LogType`, and `Summary`).

- **Utility:** `inspectiontaskbase.NewLogIngesterTask`.

### Step 4: Log Grouping & Timeline Mapping Tasks

- **Log Grouper Task:** Groups logs by a key (e.g., entity name, correlation ID) by calling Extractor functions on raw logs.
  - **Utility:** `inspectiontaskbase.NewLogGrouperTask`.
- **Timeline Mapping Task:** Maps the grouped logs to resource timelines as events or state revisions.
  - **Utility:** `inspectiontaskbase.NewLogToTimelineMapperTask`.

---

## 3. Step-by-Step Implementation Code Samples

Let's look at a concrete example of supporting a custom log type called `customapp` under `pkg/task/inspection/googlecloud/customapp/`.

### A. The Contract Package (`pkg/task/inspection/googlecloud/customapp/`)

#### `taskid.go`

Defines the TaskIDs and TaskReferences for the pipeline steps.

```go
package customapp

import (
 inspectiontaskbase "github.com/GoogleCloudPlatform/khi/pkg/core/inspection/taskbase"
 "github.com/GoogleCloudPlatform/khi/pkg/core/task/taskid"
 "github.com/GoogleCloudPlatform/khi/pkg/model/log"
)

const TaskIDPrefix = "customapp.khi.google.com/"

// 1. Form Task ID
var InputFilterKeywordTaskID = taskid.NewDefaultImplementationID[string](TaskIDPrefix + "input-keyword")

// 2. Log Query Task ID
var LogQueryTaskID = taskid.NewDefaultImplementationID[[]*log.Log](TaskIDPrefix + "query")

// 3. Log Ingestion Task ID
var LogIngesterTaskID = taskid.NewDefaultImplementationID[[]*log.Log](TaskIDPrefix + "log-ingester")

// 4. Log Grouper & Timeline Mapper Task IDs
var LogGrouperTaskID = taskid.NewDefaultImplementationID[inspectiontaskbase.LogGroupMap](TaskIDPrefix + "log-grouper")
var LogToTimelineMapperTaskID = taskid.NewDefaultImplementationID[struct{}](TaskIDPrefix + "timeline-mapper")
```

#### `extractor.go`

Defines the strongly-typed data structures and extraction functions.

> [!IMPORTANT]
>
> - **Package-level FieldPath declarations**: Pre-compiled `structured.FieldPath` values created by `structured.CompileFieldPath` are constant across log entries and MUST be declared as package-level variables in a `var (...)` block immediately below the `import` block. Never compile `FieldPath` inside functions or hot parsing loops.
> - **Non-pointer Return Values**: Extraction methods (`ExtractXXX`) MUST return value types (`FieldSet`), not pointers (`*FieldSet`). Returning values eliminates heap allocation overhead when extractors are called millions of times across high-volume log streams.
> - **Mock Support**: Extraction functions MUST check `structured.GetMock[FieldSetType](reader)` at the top of the function to allow unit tests to override extraction via `testlog.NewMockLog` / `structured.NewMockNode`.

##### Pattern 1: Direct Extractor Pattern (Single Log Source)

Used when the log format is fixed to a single ingest format (e.g., GKE Autoscaler, serial port, K8s control plane).

```go
package customapp

import (
 "github.com/GoogleCloudPlatform/khi/pkg/common/structured"
)

var (
 pathAppName   = structured.CompileFieldPath("app_name")
 pathRequestID = structured.CompileFieldPath("request_id")
 pathPayload   = structured.CompileFieldPath("payload")
)

// CustomAppFieldSet holds structured log data extracted from the log entry.
type CustomAppFieldSet struct {
 AppName   string
 RequestID string
 Payload   string
}

// ExtractCustomApp extracts CustomAppFieldSet from a raw log node reader.
func ExtractCustomApp(reader *structured.NodeReader) (CustomAppFieldSet, error) {
 if mock, ok := structured.GetMock[CustomAppFieldSet](reader); ok {
  return mock, nil
 }
 return CustomAppFieldSet{
  AppName:   reader.ReadStringOrDefault(pathAppName, "unknown-app"),
  RequestID: reader.ReadStringOrDefault(pathRequestID, ""),
  Payload:   reader.ReadStringOrDefault(pathPayload, ""),
 }, nil
}
```

##### Pattern 2: Injected Extractor Pattern (Multi-source Ingestion)

Used when the same log entity can originate from different sources with distinct field layouts (for example, K8s audit logs ingested from GCP Cloud Logging vs. OSS Kubernetes JSONL files).

The common contract defines an extractor function type and a wrapper function that retrieves the task-injected extractor from context:

```go
package k8saudit

import (
 "context"

 "github.com/GoogleCloudPlatform/khi/pkg/common/structured"
 coretask "github.com/GoogleCloudPlatform/khi/pkg/core/task"
)

// K8sAuditLogExtractor is a function type for extracting K8sAuditLogFieldSet from a NodeReader.
type K8sAuditLogExtractor func(reader *structured.NodeReader) (K8sAuditLogFieldSet, error)

// ExtractK8sAuditLog extracts K8s audit log data using the injected extractor from the task context.
func ExtractK8sAuditLog(ctx context.Context, reader *structured.NodeReader) (K8sAuditLogFieldSet, error) {
 if mock, ok := structured.GetMock[K8sAuditLogFieldSet](reader); ok {
  return mock, nil
 }
 if extractor, found := coretask.GetOptionalTaskResult(ctx, K8sAuditLogExtractorRef); found && extractor != nil {
  return extractor(reader)
 }
 return K8sAuditLogFieldSet{}, nil
}
```

#### `timeline_type.go`

Defines custom timeline types and resource verbs.

```go
package customapp

import (
 "github.com/GoogleCloudPlatform/khi/pkg/model/khifile/v6/style"
)

var (
 // TimelineTypeCustomApp is the timeline type style for Custom App resources.
 TimelineTypeCustomApp = style.MustRegisterTimelineType(
  "customapp",
  "Custom Application",
  "dns",
  0.6,
  style.ColorWhite,
  style.ColorBlack,
  style.MustForceConvertSRGBHex("#4285F4"),
  true,
  1000,
  style.AlphabeticalSortPolicy(),
 )

 // VerbCustomAppProcess is the verb style for Custom App state updates.
 VerbCustomAppProcess = style.MustRegisterVerb("Process", style.MustForceConvertSRGBHex("#0F9D58"), style.ColorWhite, true)
)
```

#### `log_type.go`

Defines custom log types.

```go
package customapp

import (
 "github.com/GoogleCloudPlatform/khi/pkg/model/khifile/v6/style"
)

var (
 // LogTypeCustomApp is the log type style for Custom App logs.
 LogTypeCustomApp = style.MustRegisterLogType(
  "customapp",
  "Custom Application Logs",
  style.MustForceConvertSRGBHex("#4285F4"),
  style.ColorWhite,
 )
)
```

#### `timeline_path.go` (Optional)

Defines helper functions to build hierarchical timeline paths.

For custom application timelines, you can define helpers to construct paths consistently. If your custom application runs as part of a Kubernetes Pod, you can build a sub-timeline path nested directly under the standard Kubernetes Pod timeline by referencing standard K8s timeline types from `inspectioncore`.

- MustXXXTimeline func must receive the context as its first argument.
- If the MustXXXTimeline func isn't for a root timeline, it must receive the parent timeline path as its second argument.

```go
package customapp

import (
 "context"

 "github.com/GoogleCloudPlatform/khi/pkg/common/khictx"
 khifilev6 "github.com/GoogleCloudPlatform/khi/pkg/model/khifile/v6"
 "github.com/GoogleCloudPlatform/khi/pkg/task/inspection/inspectioncore"
)

// MustCustomAppTimeline returns the hierarchical timeline path for a standalone Custom App.
// Constructs a path like: customapp/<appName>
func MustCustomAppTimeline(ctx context.Context, appName string) *khifilev6.TimelinePath {
 builder := khictx.MustGetValue(ctx, inspectioncore.Builder)
 return builder.TimelineAccumulator.GetPath(nil, khifilev6.PathSegment{
  Name: appName,
  Type: TimelineTypeCustomApp,
 })
}

// MustCustomAppPodTimeline returns the hierarchical timeline path for Custom App logs nested under a Pod.
// Constructs a path like: <apiVersion>/<kind>/<namespace>/<podName>/customapp
func MustCustomAppPodTimeline(ctx context.Context, podTimelinePath *khifilev6.TimelinePath) *khifilev6.TimelinePath {
  if podTimelinePath == nil || podTimelinePath.Type.GetId() != inspectioncore.TimelineTypeResource.GetId() {
  panic("parent timeline path must be Resource type")
 }

 builder := khictx.MustGetValue(ctx, inspectioncore.Builder)
 return builder.TimelineAccumulator.GetPath(podTimelinePath, khifilev6.PathSegment{
  Name: "customapp",
  Type: TimelineTypeCustomApp,
 })
}
```

---

### B. The Implementation Package (`pkg/task/inspection/googlecloud/customapp/impl/`)

#### `form.go` (Step 1)

Implements form tasks to get user-defined input.

```go
package customapp_impl

import (
 "context"

 "github.com/GoogleCloudPlatform/khi/pkg/core/inspection/formtask"
 "github.com/GoogleCloudPlatform/khi/pkg/task/inspection/googlecloud/customapp"
 "github.com/GoogleCloudPlatform/khi/pkg/task/inspection/googlecloud/gcpcommon"
)

const formPriority = gcpcommon.FormBasePriority + 5000

// InputFilterKeywordTask defines a text input form task for filtering logs.
var InputFilterKeywordTask = formtask.NewTextFormTaskBuilder(
 customapp.InputFilterKeywordTaskID,
 formPriority,
 "Filter Keyword",
).
 WithDescription("Keyword to filter Custom App logs.").
 WithDefaultValueFunc(func(ctx context.Context, previousValues []string) (string, error) {
  if len(previousValues) > 0 {
   return previousValues[0], nil
  }
  return "default-keyword", nil
 }).
 Build()
```

#### `query.go` (Step 2)

Implements querying logs from Google Cloud Logging based on parameters.

```go
package customapp_impl

import (
 "context"
 "fmt"

 coretask "github.com/GoogleCloudPlatform/khi/pkg/core/task"
 "github.com/GoogleCloudPlatform/khi/pkg/core/task/taskid"

 "github.com/GoogleCloudPlatform/khi/pkg/model/log"
 "github.com/GoogleCloudPlatform/khi/pkg/task/inspection/googlecloud/customapp"
 "github.com/GoogleCloudPlatform/khi/pkg/task/inspection/googlecloud/gcpcommon"
 "github.com/GoogleCloudPlatform/khi/pkg/task/inspection/googlecloud/k8scommon"
 "github.com/GoogleCloudPlatform/khi/pkg/task/inspection/inspectioncore"
)

// LogQueryTask executes Cloud Logging filter to fetch logs.
var LogQueryTask = gcpcommon.NewListLogEntriesTask(&customAppLogQueryTaskSetting{})

type customAppLogQueryTaskSetting struct{}

func (s *customAppLogQueryTaskSetting) TaskID() taskid.TaskImplementationID[[]*log.Log] {
 return customapp.LogQueryTaskID
}

func (s *customAppLogQueryTaskSetting) Dependencies() []taskid.UntypedTaskReference {
 return []taskid.UntypedTaskReference{
  k8scommon.ClusterIdentityTaskID.Ref(),
  customapp.InputFilterKeywordTaskID.Ref(),
 }
}

func (s *customAppLogQueryTaskSetting) Description() *gcpcommon.ListLogEntriesTaskDescription {
 return &gcpcommon.ListLogEntriesTaskDescription{

  QueryName:      "Custom App logs",
  ExampleQuery:   `resource.type="gke_cluster" AND log_id("custom-app")`,
 }
}

func (s *customAppLogQueryTaskSetting) LogFilters(ctx context.Context, taskMode inspectioncore.InspectionTaskModeType) ([]string, error) {
 keyword := coretask.GetTaskResult(ctx, customapp.InputFilterKeywordTaskID.Ref())
 query := fmt.Sprintf(`resource.type="gke_cluster" AND log_id("custom-app") AND textPayload:"%s"`, keyword)
 return []string{query}, nil
}

func (s *customAppLogQueryTaskSetting) DefaultResourceNames(ctx context.Context) ([]string, error) {
 clusterIdentity := coretask.GetTaskResult(ctx, k8scommon.ClusterIdentityTaskID.Ref())
 return []string{fmt.Sprintf("projects/%s", clusterIdentity.ProjectID)}, nil
}

func (s *customAppLogQueryTaskSetting) TimePartitionCount(ctx context.Context) (int, error) {
 return 5, nil
}

var _ gcpcommon.ListLogEntriesTaskSetting = (*customAppLogQueryTaskSetting)(nil)
```

#### `mapper.go` (Steps 3, 4)

Defines log ingestion, log grouping, and timeline mapping.

```go
package customapp_impl

import (
 "context"
 "fmt"

 "github.com/GoogleCloudPlatform/khi/pkg/common/khictx"
 inspectiontaskbase "github.com/GoogleCloudPlatform/khi/pkg/core/inspection/taskbase"
 "github.com/GoogleCloudPlatform/khi/pkg/core/task/taskid"
 khifilev6 "github.com/GoogleCloudPlatform/khi/pkg/model/khifile/v6"

 "github.com/GoogleCloudPlatform/khi/pkg/model/log"
 "github.com/GoogleCloudPlatform/khi/pkg/task/inspection/googlecloud/customapp"
 "github.com/GoogleCloudPlatform/khi/pkg/task/inspection/inspectioncore"
)

// CustomAppLogIngester V2 LogIngester (Step 3).
type CustomAppLogIngester struct{}

func (i *CustomAppLogIngester) RawLogTask() taskid.TaskReference[[]*log.Log] {
 return customapp.LogQueryTaskID.Ref()
}

func (i *CustomAppLogIngester) Dependencies() []taskid.UntypedTaskReference {
 return []taskid.UntypedTaskReference{}
}

func (i *CustomAppLogIngester) ProcessLog(ctx context.Context, l *log.Log) (*khifilev6.LogChangeSet, error) {
 cs, err := khifilev6.NewLogChangeSet(l)
 if err != nil {
  return nil, err
 }

 cs.SetLogType(customapp.LogTypeCustomApp)
 // Usually l.Timestamp from ingestion is used. However, if the log contains its own
 // custom payload field with a more precise timestamp, extract and set that instead.
 cs.SetTimestamp(l.Timestamp)

 // Extract custom fields to generate summary.
 if customFS, err := customapp.ExtractCustomApp(l.NodeReader); err == nil {
  cs.SetSummary(fmt.Sprintf("[%s] %s", customFS.AppName, customFS.Payload))
 }

 return cs, nil
}

var LogIngesterTask = inspectiontaskbase.NewLogIngesterTask(
 customapp.LogIngesterTaskID,
 &CustomAppLogIngester{},
)

// LogGrouperTask groups logs by AppName (helper for Step 4).
var LogGrouperTask = inspectiontaskbase.NewLogGrouperTask(
 customapp.LogGrouperTaskID,
 customapp.LogQueryTaskID.Ref(),
 func(ctx context.Context, l *log.Log) string {
  if customFS, err := customapp.ExtractCustomApp(l.NodeReader); err == nil {
   return customFS.AppName
  }
  return "unknown-app"
 },
)

// CustomAppTimelineMapper maps logs to timeline (Step 4).
type CustomAppTimelineMapper struct {
 inspectiontaskbase.StatelessMapperBase // Embed stateless helper.
}

func (m *CustomAppTimelineMapper) LogIngesterTask() taskid.TaskReference[[]*log.Log] {
 return customapp.LogIngesterTaskID.Ref()
}

func (m *CustomAppTimelineMapper) Dependencies() []taskid.UntypedTaskReference {
 return []taskid.UntypedTaskReference{}
}

func (m *CustomAppTimelineMapper) GroupedLogTask() taskid.TaskReference[inspectiontaskbase.LogGroupMap] {
 return customapp.LogGrouperTaskID.Ref()
}

func (m *CustomAppTimelineMapper) ProcessLogByGroup(ctx context.Context, l *log.Log, _ struct{}) (*khifilev6.TimelineChangeSet, struct{}, error) {
 customFS, err := customapp.ExtractCustomApp(l.NodeReader)
 if err != nil {
  return nil, struct{}{}, err
 }

 builder := khictx.MustGetValue(ctx, inspectioncore.CurrentV6Builder)
 targetPath := builder.TimelineAccumulator.GetPath(nil, khifilev6.PathSegment{
  Name: customFS.AppName,
  Type: customapp.TimelineTypeCustomApp,
 })

 cs := khifilev6.NewTimelineChangeSet(l)

 // Record a revision on timeline for state change.
 cs.AddRevision(targetPath, &khifilev6.StagingRevision{
  ChangedTime:  l.Timestamp,
  ResourceBody: customFS.Payload,
  VerbType:     customapp.VerbCustomAppProcess,
 })

 return cs, struct{}{}, nil
}

var LogToTimelineMapperTask = inspectiontaskbase.NewLogToTimelineMapperTask(
 customapp.LogToTimelineMapperTaskID,
 &CustomAppTimelineMapper{},
 inspectioncore.FeatureTaskLabel(
  "Custom App Logs",
  "Parser and timeline mapping for Custom App logs.",
  9000,
  false,
 ),
)

var _ inspectiontaskbase.LogToTimelineMapper[struct{}] = (*CustomAppTimelineMapper)(nil)
```

### C. Specialized Pattern: `ManifestLogToTimelineMapper` (Multi-Group Merge Mapper)

For advanced scenarios requiring the tracking and synchronization of **multiple related resource logs** chronologically (such as a parent `Pod` and its subresources like `Status` or `Binding`), KHI provides `NewManifestLogToTimelineMapper`.

This mapper automatically merges logs from multiple roles into a single stream sorted strictly by timestamp, and passes the state `T` across all events.

#### Key Interfaces and Structures

- **`RelatedGroupSet`**: Groups related logs by role name (e.g., `"source" -> PodGroup`, `"target" -> BindingGroup`).
- **`MultiGroupLogEvent`**: Contains the currently yielding `Log`, the role (`GroupRole`), and the helper methods:
  - `GetLastBodyReader(role string) (*structured.NodeReader, bool)`: Retrieves the latest manifest body of the specified role as a `NodeReader` at the time of the event using highly optimized `O(log N)` binary search.
  - `GetLastBodyYAML(role string) (string, bool)`: Retrieves the latest manifest body as a YAML string.

#### Code Sample: Single-Pass Stateful Manifest Mapper

```go
package myapp_impl

import (
 "context"
 "time"

 "github.com/GoogleCloudPlatform/khi/pkg/core/task/taskid"
 pb "github.com/GoogleCloudPlatform/khi/pkg/generated/khifile/v6"
 khifilev6 "github.com/GoogleCloudPlatform/khi/pkg/model/khifile/v6"
 "github.com/GoogleCloudPlatform/khi/pkg/model/log"
 "github.com/GoogleCloudPlatform/khi/pkg/task/inspection/common/k8saudit"
)

type MyState struct {
 WasDeleted bool
}

type MyManifestMapper struct {
 // Embeds single pass helper.
 k8saudit.ManifestSinglePassMapperBase[*MyState]
}

func (m *MyManifestMapper) TaskID() taskid.TaskImplementationID[struct{}] {
 return mycontract.MyManifestMapperTaskID
}

func (m *MyManifestMapper) LogIngesterTask() taskid.TaskReference[[]*log.Log] {
 return k8saudit.K8sAuditLogIngesterTaskID.Ref()
}

func (m *MyManifestMapper) GroupedLogTask() taskid.TaskReference[k8saudit.ResourceManifestLogGroupMap] {
 return k8saudit.ResourceLifetimeTrackerTaskID.Ref()
}

func (m *MyManifestMapper) Dependencies() []taskid.UntypedTaskReference {
 return []taskid.UntypedTaskReference{}
}

// ResolveRelatedGroupSets groups a parent resource (source) and its subresource (target) together.
func (m *MyManifestMapper) ResolveRelatedGroupSets(ctx context.Context, groupedLogs k8saudit.ResourceManifestLogGroupMap) ([]k8saudit.RelatedGroupSet, error) {
 result := []k8saudit.RelatedGroupSet{}
 for _, group := range groupedLogs {
  if group.Resource.Type() == k8saudit.Subresource {
   parentGroup := groupedLogs[group.Resource.ParentIdentity().ResourcePathString()]
   result = append(result, k8saudit.RelatedGroupSet{
    Roles: map[string]*k8saudit.ResourceManifestLogGroup{
     "source": parentGroup,
     "target": group,
    },
   })
  }
 }
 return result, nil
}

// ProcessLog processes chronologically merged events.
func (m *MyManifestMapper) ProcessLog(ctx context.Context, event k8saudit.MultiGroupLogEvent, state *MyState) (*khifilev6.TimelineChangeSet, *MyState, error) {
 if state == nil {
  state = &MyState{}
 }

 cs := khifilev6.NewTimelineChangeSet(event.Log)

 // Handle parent deletion event to propagate deletion to the subresource.
 if event.GroupRole == "source" && event.EventType == k8saudit.ChangeEventTypeDeletion {
  targetGroup := event.GroupSet.Roles["target"]
  targetPath := MustResolveTimelinePath(ctx, targetGroup.Resource)

  cs.AddRevision(targetPath, &khifilev6.StagingRevision{
   ChangedTime: time.Now(),
   StateType:   k8saudit.RevisionStateK8sResourceIsDeleted,
  })
  state.WasDeleted = true
 }

 return cs, state, nil
}

var _ k8saudit.ManifestLogToTimelineMapper[*MyState] = (*MyManifestMapper)(nil)
```

#### `registration.go`

Registers the tasks with the central registry.

```go
package customapp_impl

import (
 coreinspection "github.com/GoogleCloudPlatform/khi/pkg/core/inspection"
 coretask "github.com/GoogleCloudPlatform/khi/pkg/core/task"
)

// Register registers all customapp tasks to the central registry.
func Register(registry coreinspection.InspectionTaskRegistry) error {
 return coretask.RegisterTasks(
  registry,
  InputFilterKeywordTask,
  LogQueryTask,
  LogIngesterTask,
  LogGrouperTask,
  LogToTimelineMapperTask,
 )
}
```

---

## 4. Testing Log Parsers

Refer to [log-timeline-mapper](skill://log-timeline-mapper) for detailed unit testing strategies of `LogIngester` and `LogToTimelineMapper`.

### Testing `ManifestLogToTimelineMapper`

Since `ManifestLogToTimelineMapper` coordinates chronologically merged streams and tracks previous states, testing it requires:

1. **Chronological Merge Validation**: Testing that events from different roles are merged correctly.
2. **Historical Snapshot Validation**: Testing that `GetLastBodyReader` or `GetLastBodyYAML` accurately yields the snapshot of other roles at the event's timestamp.

Use a **Table-Driven Test** pattern to verify these behaviors comprehensively.

#### Test Example: Table-Driven Snapshot Verification

```go
func TestGetLastBody(t *testing.T) {
 t1 := time.Date(2026, 5, 26, 10, 0, 0, 0, time.UTC)
 t2 := t1.Add(time.Minute)

 nodeA1, _ := structured.FromGoValue(map[string]any{"value": "A1"}, &structured.AlphabeticalGoMapKeyOrderProvider{})
 nodeB1, _ := structured.FromGoValue(map[string]any{"value": "B1"}, &structured.AlphabeticalGoMapKeyOrderProvider{})

 logA1 := testlog.NewMockLog(t1)
 logB1 := testlog.NewMockLog(t2)

 groupSet := RelatedGroupSet{
  Roles: map[string]*ResourceManifestLogGroup{
   "roleA": {
    Logs: []*ResourceManifestLog{
     {Log: logA1, ResourceBodyYAML: "value: A1", ResourceBodyReader: structured.NewNodeReader(nodeA1)},
    },
   },
   "roleB": {
    Logs: []*ResourceManifestLog{
     {Log: logB1, ResourceBodyYAML: "value: B1", ResourceBodyReader: structured.NewNodeReader(nodeB1)},
    },
   },
  },
 }

 events := make([]MultiGroupLogEvent, 0)
 for event := range iterateMultiGroupLog(groupSet) {
  events = append(events, event)
 }

 testCases := []struct {
  name         string
  eventIndex   int
  expectedRole string
  roleToCheck  string
  wantFound    bool
  wantYAML     string
 }{
  {
   name:         "event 0: check roleA body",
   eventIndex:   0,
   expectedRole: "roleA",
   roleToCheck:  "roleA",
   wantFound:    true,
   wantYAML:     "value: A1",
  },
  {
   name:         "event 0: check roleB body (not exist yet)",
   eventIndex:   0,
   expectedRole: "roleA",
   roleToCheck:  "roleB",
   wantFound:    false,
  },
  {
   name:         "event 1: check roleA body from roleB event",
   eventIndex:   1,
   expectedRole: "roleB",
   roleToCheck:  "roleA",
   wantFound:    true,
   wantYAML:     "value: A1",
  },
 }

 for _, tc := range testCases {
  t.Run(tc.name, func(t *testing.T) {
   e := events[tc.eventIndex]

   if e.GroupRole != tc.expectedRole {
    t.Errorf("expected group role %q, got %q", tc.expectedRole, e.GroupRole)
   }

   yaml, ok := e.GetLastBodyYAML(tc.roleToCheck)
   if ok != tc.wantFound {
    t.Errorf("GetLastBodyYAML(%q) ok = %t, want %t", tc.roleToCheck, ok, tc.wantFound)
   }
   if ok && yaml != tc.wantYAML {
    t.Errorf("GetLastBodyYAML(%q) = %q, want %q", tc.roleToCheck, yaml, tc.wantYAML)
   }
  })
 }
}
```

### Testing Tasks Implemented with `ManifestLogToTimelineMapper`

To unit test a concrete mapper task implementing `ManifestLogToTimelineMapper[T]`, you should isolate and test its `ProcessLog` (and `PreProcessLog`) method using table-driven tests.

The test setup requires:

1. **v6 Builder Initialization**: Instantiate a `khifilev6.Builder` and construct the expected `TimelinePath` instances.
2. **Context Injection**: Inject the builder into the test context utilizing `khictx.WithValue` and the key `inspectioncore.Builder`.
3. **Mock Event Construction**: Manually instantiate a `MultiGroupLogEvent` with mock logs and roles, and supply a mock `RelatedGroupSet` if testing body-reference lookups.
4. **Fluent ChangeSet Assertions**: Verify the generated timelines using the fluent asserter utility `testchangeset.AssertTimeline`.

#### Task Unit Test Example

This example isolates and tests the `MyManifestMapper` defined in Section 3.C.

```go
func TestMyManifestMapper_ProcessLog(t *testing.T) {
 // 1. Set up the mock Builder and construct comparison paths hierarchically.
 builder := khifilev6.NewBuilder()
 cluster := builder.TimelineAccumulator.GetPath(nil, khifilev6.PathSegment{Name: "k8s", Type: inspectioncore.TimelineTypeK8sCluster})
 api := builder.TimelineAccumulator.GetPath(cluster, khifilev6.PathSegment{Name: "core/v1", Type: inspectioncore.TimelineTypeAPIVersion})
 kind := builder.TimelineAccumulator.GetPath(api, khifilev6.PathSegment{Name: "pod", Type: inspectioncore.TimelineTypeKind})
 ns := builder.TimelineAccumulator.GetPath(kind, khifilev6.PathSegment{Name: "default", Type: inspectioncore.TimelineTypeNamespace})
 pod := builder.TimelineAccumulator.GetPath(ns, khifilev6.PathSegment{Name: "my-pod", Type: inspectioncore.TimelineTypeResource})
 targetPath := builder.TimelineAccumulator.GetPath(pod, khifilev6.PathSegment{Name: "binding", Type: TimelineTypeSubresource})

 testCases := []struct {
  name      string
  event     MultiGroupLogEvent
  prevState *MyState
  assert    func(t *testing.T, ctx context.Context, cs *khifilev6.TimelineChangeSet, state *MyState)
 }{
  {
   name: "parent pod deletion propagates delete revision to subresource binding",
   event: MultiGroupLogEvent{
    Log: testlog.NewMockLog(time.Date(2026, 5, 26, 10, 0, 0, 0, time.UTC)),
    GroupRole: "source", // Parent resource
    EventType: ChangeEventTypeDeletion,
    GroupSet: RelatedGroupSet{
     Roles: map[string]*ResourceManifestLogGroup{
      "target": {
       Resource: &ResourceIdentity{
        APIVersion:      "core/v1",
        Kind:            "pod",
        Name:            "my-pod",
        Namespace:       "default",
        SubresourceName: "binding",
       },
      },
     },
    },
   },
   prevState: &MyState{WasDeleted: false},
   assert: func(t *testing.T, ctx context.Context, cs *khifilev6.TimelineChangeSet, state *MyState) {
    // Verify the generated staged revisions using fluent assertions.
    testchangeset.AssertTimeline(t, cs).
     HasRevision(targetPath, &khifilev6.StagingRevision{
      StateType: k8saudit.RevisionStateK8sResourceIsDeleted,
     })

    // Verify that state changes are correctly tracked.
    if !state.WasDeleted {
     t.Errorf("state.WasDeleted = false, want true")
    }
   },
  },
 }

 mapper := &MyManifestMapper{}
 for _, tc := range testCases {
  t.Run(tc.name, func(t *testing.T) {
   // 2. Inject the SAME Builder instance into context.
   ctx := khictx.WithValue(t.Context(), inspectioncore.Builder, builder)

   // 3. Call ProcessLog directly.
   cs, nextState, err := mapper.ProcessLog(ctx, tc.event, tc.prevState)
   if err != nil {
    t.Fatalf("ProcessLog() returned unexpected error: %v", err)
   }

   // 4. Assert outcomes.
   tc.assert(t, ctx, cs, nextState)
  })
 }
}
```

Discussion

Did this work in your project? Say what you used it for and what you changed. People and their agents can both post here.

Posts are public.Sign in to post

No one has posted yet. Be the first.