feat(go): add Phase 1 parity sidecar
This commit is contained in:
parent
65f5be6420
commit
f02509aab5
22 changed files with 947 additions and 6 deletions
119
go-agent/internal/agent/agent.go
Normal file
119
go-agent/internal/agent/agent.go
Normal file
|
|
@ -0,0 +1,119 @@
|
|||
// SPDX-License-Identifier: GPL-3.0-or-later
|
||||
|
||||
// Package agent owns the parallel Go sweep loop. It emits JSON-compatible
|
||||
// records but does not write Python daemon state or replace the production
|
||||
// service during Phase 1.
|
||||
package agent
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"os"
|
||||
"time"
|
||||
|
||||
"codeberg.org/anassaeneroi/enodia-sentinal/go-agent/internal/config"
|
||||
"codeberg.org/anassaeneroi/enodia-sentinal/go-agent/internal/detectors"
|
||||
"codeberg.org/anassaeneroi/enodia-sentinal/go-agent/internal/model"
|
||||
"codeberg.org/anassaeneroi/enodia-sentinal/go-agent/internal/schema"
|
||||
)
|
||||
|
||||
const Version = "0.7.0"
|
||||
|
||||
// CaptureFunc returns one injectable SystemState-equivalent sweep.
|
||||
type CaptureFunc func() (model.State, error)
|
||||
|
||||
// EmitFunc receives one enodia.event.v1-compatible record.
|
||||
type EmitFunc func(map[string]any) error
|
||||
|
||||
// Agent is the Phase 1 sidecar sweep loop.
|
||||
type Agent struct {
|
||||
Config config.Config
|
||||
Capture CaptureFunc
|
||||
Host func() (string, error)
|
||||
Now func() time.Time
|
||||
}
|
||||
|
||||
// New builds an agent with production clock and hostname providers.
|
||||
func New(cfg config.Config, capture CaptureFunc) *Agent {
|
||||
return &Agent{
|
||||
Config: cfg,
|
||||
Capture: capture,
|
||||
Host: os.Hostname,
|
||||
Now: time.Now,
|
||||
}
|
||||
}
|
||||
|
||||
// Sweep captures state once and returns alert events followed by one status
|
||||
// event. Retained snapshot fields remain empty until persistence is ported.
|
||||
func (a *Agent) Sweep() ([]map[string]any, error) {
|
||||
state, err := a.Capture()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
host, err := a.Host()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
timestamp := a.Now().Local().Format(time.RFC3339Nano)
|
||||
|
||||
alerts := make([]model.Alert, 0)
|
||||
if a.Config.Enabled("deleted_exe") {
|
||||
alerts = append(alerts, detectors.DeletedExe(state)...)
|
||||
}
|
||||
events := make([]map[string]any, 0, len(alerts)+1)
|
||||
for _, alert := range alerts {
|
||||
record, err := schema.Build("alert", alert, host, timestamp)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
events = append(events, record)
|
||||
}
|
||||
status := schema.Status{
|
||||
Schema: schema.StatusV1,
|
||||
Version: Version,
|
||||
Running: true,
|
||||
TotalAlerts: 0,
|
||||
Counts: map[string]int{},
|
||||
LastAlert: nil,
|
||||
EBPF: "unknown",
|
||||
EBPFExec: "unknown",
|
||||
EBPFSyscall: "unknown",
|
||||
Host: host,
|
||||
}
|
||||
record, err := schema.Build("status", status, host, timestamp)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return append(events, record), nil
|
||||
}
|
||||
|
||||
// Run emits one sweep immediately and then waits for each configured interval.
|
||||
// once is used by parity checks and operator-visible smoke tests.
|
||||
func (a *Agent) Run(ctx context.Context, once bool, emit EmitFunc) error {
|
||||
if a.Capture == nil || emit == nil {
|
||||
return fmt.Errorf("capture and emit functions are required")
|
||||
}
|
||||
for {
|
||||
events, err := a.Sweep()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
for _, event := range events {
|
||||
if err := emit(event); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
if once {
|
||||
return nil
|
||||
}
|
||||
timer := time.NewTimer(a.Config.SampleInterval)
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
if !timer.Stop() {
|
||||
<-timer.C
|
||||
}
|
||||
return nil
|
||||
case <-timer.C:
|
||||
}
|
||||
}
|
||||
}
|
||||
54
go-agent/internal/agent/agent_test.go
Normal file
54
go-agent/internal/agent/agent_test.go
Normal file
|
|
@ -0,0 +1,54 @@
|
|||
// SPDX-License-Identifier: GPL-3.0-or-later
|
||||
|
||||
package agent
|
||||
|
||||
import (
|
||||
"context"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"codeberg.org/anassaeneroi/enodia-sentinal/go-agent/internal/config"
|
||||
"codeberg.org/anassaeneroi/enodia-sentinal/go-agent/internal/model"
|
||||
)
|
||||
|
||||
func TestRunOnceEmitsAlertAndStatusEnvelopes(t *testing.T) {
|
||||
agent := New(config.Default(), func() (model.State, error) {
|
||||
return model.State{Processes: []model.Process{
|
||||
{PID: 42, Comm: "dropper", Exe: "/tmp/dropper (deleted)"},
|
||||
}}, nil
|
||||
})
|
||||
agent.Host = func() (string, error) { return "host-a", nil }
|
||||
agent.Now = func() time.Time {
|
||||
return time.Date(2026, 7, 10, 0, 0, 0, 0, time.FixedZone("PDT", -7*3600))
|
||||
}
|
||||
var events []map[string]any
|
||||
err := agent.Run(context.Background(), true, func(event map[string]any) error {
|
||||
events = append(events, event)
|
||||
return nil
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if len(events) != 2 || events[0]["event_type"] != "alert" || events[1]["event_type"] != "status" {
|
||||
t.Fatalf("unexpected events: %#v", events)
|
||||
}
|
||||
alert := events[0]["alert"].(model.Alert)
|
||||
if alert.SID != 100012 || alert.Detail != "pid=42 comm=dropper exe=[/tmp/dropper (deleted)]" {
|
||||
t.Fatalf("unexpected alert payload: %#v", alert)
|
||||
}
|
||||
}
|
||||
|
||||
func TestDisabledDetectorEmitsStatusOnly(t *testing.T) {
|
||||
cfg := config.Default()
|
||||
cfg.Detectors = map[string]bool{}
|
||||
agent := New(cfg, func() (model.State, error) {
|
||||
return model.State{Processes: []model.Process{{PID: 42, Exe: "/tmp/x (deleted)"}}}, nil
|
||||
})
|
||||
events, err := agent.Sweep()
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if len(events) != 1 || events[0]["event_type"] != "status" {
|
||||
t.Fatalf("unexpected events: %#v", events)
|
||||
}
|
||||
}
|
||||
Loading…
Add table
Add a link
Reference in a new issue