Watch
1
0
Fork
You've already forked RedFlag
0

WIP: Vanguard uncommitted changes - aggregator loop, client, handlers, lifecycle

This commit is contained in:
Ani Tunturi 2026-04-13 17:51:48 -04:00
commit 02bbc8d2c8
10 changed files with 351 additions and 18 deletions

View file

@ -344,6 +344,9 @@ func processCommands(ctx *loopContext, commands []client.Command) {
handlers.HandleScanWindows(ctx.apiClient, ctx.cfg, ctx.ackTracker, ctx.scanOrchestrator, cmd.ID)
case "scan_winget":
handlers.HandleScanWinget(ctx.apiClient, ctx.cfg, ctx.ackTracker, ctx.scanOrchestrator, cmd.ID)
case "scan_updates":
// Virtual subsystem - triggers all package scanners
handlers.HandleScanUpdates(ctx.apiClient, ctx.cfg, ctx.ackTracker, ctx.scanOrchestrator, cmd.ID)
default:
// TODO: Move other command handlers from main.go to handlers package
log.Printf("Command type %s not yet refactored to handlers package", cmd.Type)

View file

@ -184,6 +184,7 @@ type RegisterRequest struct {
MachineID string `json:"machine_id"`
PublicKeyFingerprint string `json:"public_key_fingerprint"`
Metadata map[string]string `json:"metadata"`
AvailableScanners []string `json:"available_scanners"` // Platform-specific package managers (apt, dnf, winget, windows)
}
// RegisterResponse is returned after successful registration

View file

@ -4,6 +4,8 @@ import (
"context"
"fmt"
"log"
"runtime"
"strings"
"time"
"github.com/Fimeg/RedFlag/aggregator-agent/internal/acknowledgment"
@ -11,6 +13,7 @@ import (
"github.com/Fimeg/RedFlag/aggregator-agent/internal/config"
"github.com/Fimeg/RedFlag/aggregator-agent/internal/models"
"github.com/Fimeg/RedFlag/aggregator-agent/internal/orchestrator"
"github.com/Fimeg/RedFlag/aggregator-agent/internal/scanner"
)
// reportLogWithAck reports a command log to the server and tracks it for acknowledgment
@ -33,6 +36,160 @@ func reportLogWithAck(apiClient *client.Client, cfg *config.Config, ackTracker *
return nil
}
// HandleScanUpdates scans for ALL package updates across all available package managers
// This is the virtual "updates" subsystem that triggers apt, dnf, winget, and windows scans
func HandleScanUpdates(apiClient *client.Client, cfg *config.Config, ackTracker *acknowledgment.Tracker, orch *orchestrator.Orchestrator, commandID string) error {
log.Println("Scanning for all package updates...")
ctx := context.Background()
startTime := time.Now()
var totalUpdates int
var scanResults []string
var errors []string
// Detect OS and run appropriate scanners
osType := runtime.GOOS
// Linux: Try APT first, then DNF
if osType == "linux" {
// Try APT
aptScanner := orchestrator.NewAPTScannerWrapper(scanner.NewAPTScanner())
if aptScanner.IsAvailable() {
log.Println("[updates] Running APT scan...")
result, err := orch.ScanSingle(ctx, "apt")
if err != nil {
errors = append(errors, fmt.Sprintf("APT: %v", err))
} else {
scanResults = append(scanResults, fmt.Sprintf("APT: %d updates", len(result.Updates)))
totalUpdates += len(result.Updates)
// Report APT updates
if len(result.Updates) > 0 {
report := client.UpdateReport{
CommandID: commandID,
Timestamp: time.Now(),
Updates: result.Updates,
}
if err := apiClient.ReportUpdates(cfg.AgentID, report); err != nil {
log.Printf("[WARNING] [updates] Failed to report APT updates: %v", err)
}
}
}
}
// Try DNF
dnfScanner := orchestrator.NewDNFScannerWrapper(scanner.NewDNFScanner())
if dnfScanner.IsAvailable() {
log.Println("[updates] Running DNF scan...")
result, err := orch.ScanSingle(ctx, "dnf")
if err != nil {
errors = append(errors, fmt.Sprintf("DNF: %v", err))
} else {
scanResults = append(scanResults, fmt.Sprintf("DNF: %d updates", len(result.Updates)))
totalUpdates += len(result.Updates)
// Report DNF updates
if len(result.Updates) > 0 {
report := client.UpdateReport{
CommandID: commandID,
Timestamp: time.Now(),
Updates: result.Updates,
}
if err := apiClient.ReportUpdates(cfg.AgentID, report); err != nil {
log.Printf("[WARNING] [updates] Failed to report DNF updates: %v", err)
}
}
}
}
}
// Windows: Try Windows Update and Winget
if osType == "windows" {
// Try Windows Update
windowsScanner := orchestrator.NewWindowsUpdateScannerWrapper(scanner.NewWindowsUpdateScanner())
if windowsScanner.IsAvailable() {
log.Println("[updates] Running Windows Update scan...")
result, err := orch.ScanSingle(ctx, "windows")
if err != nil {
errors = append(errors, fmt.Sprintf("Windows: %v", err))
} else {
scanResults = append(scanResults, fmt.Sprintf("Windows: %d updates", len(result.Updates)))
totalUpdates += len(result.Updates)
// Report Windows updates
if len(result.Updates) > 0 {
report := client.UpdateReport{
CommandID: commandID,
Timestamp: time.Now(),
Updates: result.Updates,
}
if err := apiClient.ReportUpdates(cfg.AgentID, report); err != nil {
log.Printf("[WARNING] [updates] Failed to report Windows updates: %v", err)
}
}
}
}
// Try Winget
wingetScanner := orchestrator.NewWingetScannerWrapper(scanner.NewWingetScanner())
if wingetScanner.IsAvailable() {
log.Println("[updates] Running Winget scan...")
result, err := orch.ScanSingle(ctx, "winget")
if err != nil {
errors = append(errors, fmt.Sprintf("Winget: %v", err))
} else {
scanResults = append(scanResults, fmt.Sprintf("Winget: %d updates", len(result.Updates)))
totalUpdates += len(result.Updates)
// Report Winget updates
if len(result.Updates) > 0 {
report := client.UpdateReport{
CommandID: commandID,
Timestamp: time.Now(),
Updates: result.Updates,
}
if err := apiClient.ReportUpdates(cfg.AgentID, report); err != nil {
log.Printf("[WARNING] [updates] Failed to report Winget updates: %v", err)
}
}
}
}
}
duration := time.Since(startTime)
stdout := fmt.Sprintf("Package update scan completed in %.2f seconds\n\nResults:\n%s\n\nTotal updates found: %d",
duration.Seconds(),
strings.Join(scanResults, "\n"),
totalUpdates)
stderr := ""
exitCode := 0
if len(errors) > 0 {
stderr = fmt.Sprintf("Errors encountered:\n%s", strings.Join(errors, "\n"))
exitCode = 1
}
// Create history entry
logReport := client.LogReport{
CommandID: commandID,
Action: "scan_updates",
Result: map[bool]string{true: "success", false: "partial_failure"}[exitCode == 0],
Stdout: stdout,
Stderr: stderr,
ExitCode: exitCode,
DurationSeconds: int(duration.Seconds()),
Metadata: map[string]string{
"subsystem_label": "Package Updates",
"subsystem": "updates",
"total_updates": fmt.Sprintf("%d", totalUpdates),
"scanners_run": fmt.Sprintf("%d", len(scanResults)),
},
}
if err := reportLogWithAck(apiClient, cfg, ackTracker, logReport); err != nil {
log.Printf("[ERROR] [agent] [updates] report_log_failed: %v", err)
} else {
log.Printf("[INFO] [agent] [updates] scan completed: %d updates found", totalUpdates)
}
return nil
}
// HandleScanStorage scans disk usage metrics only
func HandleScanStorage(apiClient *client.Client, cfg *config.Config, ackTracker *acknowledgment.Tracker, orch *orchestrator.Orchestrator, commandID string) error {
log.Println("Scanning storage...")

View file

@ -10,6 +10,7 @@ import (
"github.com/Fimeg/RedFlag/aggregator-agent/internal/config"
"github.com/Fimeg/RedFlag/aggregator-agent/internal/constants"
"github.com/Fimeg/RedFlag/aggregator-agent/internal/crypto"
"github.com/Fimeg/RedFlag/aggregator-agent/internal/scanner"
"github.com/Fimeg/RedFlag/aggregator-agent/internal/system"
"github.com/Fimeg/RedFlag/aggregator-agent/internal/version"
)
@ -80,6 +81,10 @@ func RegisterAgent(cfg *config.Config, serverURL string) error {
log.Printf("Warning: No embedded public key fingerprint found")
}
// Detect available scanners for platform-specific subsystem creation
availableScanners := detectAvailableScanners()
log.Printf("[INFO] [agent] [registration] detected_scanners=%v", availableScanners)
req := client.RegisterRequest{
Hostname: sysInfo.Hostname,
OSType: sysInfo.OSType,
@ -89,6 +94,7 @@ func RegisterAgent(cfg *config.Config, serverURL string) error {
MachineID: machineID,
PublicKeyFingerprint: publicKeyFingerprint,
Metadata: metadata,
AvailableScanners: availableScanners,
}
resp, err := apiClient.Register(req)
@ -132,3 +138,36 @@ func FetchAndCachePublicKey(serverURL string) error {
_, err := crypto.FetchAndCacheServerPublicKey(serverURL)
return err
}
// detectAvailableScanners checks which package manager scanners are available on this system
func detectAvailableScanners() []string {
scanners := []string{}
// Check each scanner's availability using the scanner package directly
aptScanner := scanner.NewAPTScanner()
if aptScanner.IsAvailable() {
scanners = append(scanners, "apt")
}
dnfScanner := scanner.NewDNFScanner()
if dnfScanner.IsAvailable() {
scanners = append(scanners, "dnf")
}
wingetScanner := scanner.NewWingetScanner()
if wingetScanner.IsAvailable() {
scanners = append(scanners, "winget")
}
windowsScanner := scanner.NewWindowsUpdateScanner()
if windowsScanner.IsAvailable() {
scanners = append(scanners, "windows")
}
dockerScanner, _ := scanner.NewDockerScanner()
if dockerScanner != nil && dockerScanner.IsAvailable() {
scanners = append(scanners, "docker")
}
return scanners
}

View file

@ -204,6 +204,50 @@ func (h *AgentHandler) RegisterAgent(c *gin.Context) {
return
}
// Step 4: Create platform-specific subsystems based on agent's available scanners
// This replaces the generic "updates" subsystem with specific ones (apt, dnf, winget, windows)
if len(req.AvailableScanners) > 0 {
for _, scanner := range req.AvailableScanners {
intervalMinutes := 60 // Default 1 hour for update scanners
sub := models.AgentSubsystem{
AgentID: agent.ID,
Subsystem: scanner, // apt, dnf, winget, windows, docker
Enabled: true,
AutoRun: true,
IntervalMinutes: intervalMinutes,
}
if err := h.subsystemQueries.CreateSubsystemWithTx(tx, &sub); err != nil {
log.Printf("[WARNING] [server] [registration] create_subsystem_failed scanner=%s error=%v", scanner, err)
// Non-fatal - continue with other subsystems
}
}
// Always add storage, system, docker as generic subsystems
genericSubsystems := []struct {
name string
interval int
}{
{"storage", 5},
{"system", 5},
}
for _, gen := range genericSubsystems {
sub := models.AgentSubsystem{
AgentID: agent.ID,
Subsystem: gen.name,
Enabled: true,
AutoRun: true,
IntervalMinutes: gen.interval,
}
if err := h.subsystemQueries.CreateSubsystemWithTx(tx, &sub); err != nil {
log.Printf("[WARNING] [server] [registration] create_generic_subsystem_failed subsystem=%s error=%v", gen.name, err)
}
}
} else {
// Fallback: create default subsystems if agent didn't report available scanners
if err := h.subsystemQueries.CreateDefaultSubsystemsWithTx(tx, agent.ID, nil); err != nil {
log.Printf("[WARNING] [server] [registration] create_default_subsystems_failed error=%v", err)
}
}
// Commit transaction — all DB operations succeed or none do
if err := tx.Commit(); err != nil {
log.Printf("[ERROR] [server] [registration] transaction_commit_failed error=%q", err)

View file

@ -9,6 +9,16 @@ import (
"github.com/jmoiron/sqlx"
)
// SubsystemQueriesInterface defines the interface for subsystem operations (enables transaction support)
type SubsystemQueriesInterface interface {
GetSubsystems(agentID uuid.UUID) ([]models.AgentSubsystem, error)
GetSubsystem(agentID uuid.UUID, subsystem string) (*models.AgentSubsystem, error)
CreateSubsystem(sub *models.AgentSubsystem) error
CreateSubsystemWithTx(tx *sqlx.Tx, sub *models.AgentSubsystem) error
CreateDefaultSubsystems(agentID uuid.UUID, availableScanners []string) error
CreateDefaultSubsystemsWithTx(tx *sqlx.Tx, agentID uuid.UUID, availableScanners []string) error
}
type SubsystemQueries struct {
db *sqlx.DB
}
@ -232,13 +242,25 @@ func (q *SubsystemQueries) SetInterval(agentID uuid.UUID, subsystem string, inte
// CreateSubsystem creates a new subsystem configuration (used for custom subsystems)
func (q *SubsystemQueries) CreateSubsystem(sub *models.AgentSubsystem) error {
return q.createSubsystemInternal(q.db, sub)
}
// CreateSubsystemWithTx creates a subsystem within a transaction
func (q *SubsystemQueries) CreateSubsystemWithTx(tx *sqlx.Tx, sub *models.AgentSubsystem) error {
return q.createSubsystemInternal(tx, sub)
}
// createSubsystemInternal is the internal implementation that works with both DB and transactions
func (q *SubsystemQueries) createSubsystemInternal(exec interface {
QueryRow(query string, args ...interface{}) *sql.Row
}, sub *models.AgentSubsystem) error {
query := `
INSERT INTO agent_subsystems (agent_id, subsystem, enabled, interval_minutes, auto_run, last_run_at, next_run_at)
VALUES ($1, $2, $3, $4, $5, $6, $7)
RETURNING id, created_at, updated_at
`
err := q.db.QueryRow(
err := exec.QueryRow(
query,
sub.AgentID,
sub.Subsystem,
@ -281,16 +303,43 @@ func (q *SubsystemQueries) DeleteSubsystem(agentID uuid.UUID, subsystem string)
}
// CreateDefaultSubsystems creates default subsystems for a new agent
func (q *SubsystemQueries) CreateDefaultSubsystems(agentID uuid.UUID) error {
// Now creates platform-specific subsystems (apt, dnf, winget, windows) based on available scanners
func (q *SubsystemQueries) CreateDefaultSubsystems(agentID uuid.UUID, availableScanners []string) error {
return q.createDefaultSubsystemsInternal(q.db, agentID, availableScanners)
}
// CreateDefaultSubsystemsWithTx creates default subsystems within a transaction
func (q *SubsystemQueries) CreateDefaultSubsystemsWithTx(tx *sqlx.Tx, agentID uuid.UUID, availableScanners []string) error {
return q.createDefaultSubsystemsInternal(tx, agentID, availableScanners)
}
// createDefaultSubsystemsInternal is the internal implementation
func (q *SubsystemQueries) createDefaultSubsystemsInternal(exec interface {
QueryRow(query string, args ...interface{}) *sql.Row
}, agentID uuid.UUID, availableScanners []string) error {
// Always create storage, system, docker as generic subsystems
defaults := []models.AgentSubsystem{
{AgentID: agentID, Subsystem: "updates", Enabled: true, AutoRun: true, IntervalMinutes: 60},
{AgentID: agentID, Subsystem: "storage", Enabled: true, AutoRun: true, IntervalMinutes: 5},
{AgentID: agentID, Subsystem: "system", Enabled: true, AutoRun: true, IntervalMinutes: 5},
{AgentID: agentID, Subsystem: "docker", Enabled: true, AutoRun: true, IntervalMinutes: 15},
}
// Create platform-specific package manager subsystems based on available scanners
for _, scanner := range availableScanners {
switch scanner {
case "apt":
defaults = append(defaults, models.AgentSubsystem{AgentID: agentID, Subsystem: "apt", Enabled: true, AutoRun: true, IntervalMinutes: 15})
case "dnf":
defaults = append(defaults, models.AgentSubsystem{AgentID: agentID, Subsystem: "dnf", Enabled: true, AutoRun: true, IntervalMinutes: 15})
case "winget":
defaults = append(defaults, models.AgentSubsystem{AgentID: agentID, Subsystem: "winget", Enabled: true, AutoRun: true, IntervalMinutes: 60})
case "windows":
defaults = append(defaults, models.AgentSubsystem{AgentID: agentID, Subsystem: "windows", Enabled: true, AutoRun: true, IntervalMinutes: 60})
}
}
for _, sub := range defaults {
if err := q.CreateSubsystem(&sub); err != nil {
if err := q.createSubsystemInternal(exec, &sub); err != nil {
return fmt.Errorf("failed to create subsystem %s: %w", sub.Subsystem, err)
}
}

View file

@ -79,15 +79,16 @@ type AgentSpecs struct {
// AgentRegistrationRequest is the payload for agent registration
type AgentRegistrationRequest struct {
Hostname string `json:"hostname" binding:"required"`
OSType string `json:"os_type" binding:"required"`
OSVersion string `json:"os_version"`
OSArchitecture string `json:"os_architecture"`
AgentVersion string `json:"agent_version" binding:"required"`
RegistrationToken string `json:"registration_token"` // Optional, for fallback method
MachineID string `json:"machine_id"` // Unique machine identifier
PublicKeyFingerprint string `json:"public_key_fingerprint"` // Embedded public key fingerprint
Metadata map[string]string `json:"metadata"`
Hostname string `json:"hostname" binding:"required"`
OSType string `json:"os_type" binding:"required"`
OSVersion string `json:"os_version"`
OSArchitecture string `json:"os_architecture"`
AgentVersion string `json:"agent_version" binding:"required"`
RegistrationToken string `json:"registration_token"` // Optional, for fallback method
MachineID string `json:"machine_id"` // Unique machine identifier
PublicKeyFingerprint string `json:"public_key_fingerprint"` // Embedded public key fingerprint
Metadata map[string]string `json:"metadata"`
AvailableScanners []string `json:"available_scanners"` // Platform-specific scanners (apt, dnf, winget, windows, docker)
}
// AgentRegistrationResponse is returned after successful registration

View file

@ -4,6 +4,7 @@ import (
"context"
"fmt"
"log"
"strings"
"time"
"github.com/Fimeg/RedFlag/aggregator-server/internal/config"
@ -197,8 +198,12 @@ func (s *AgentLifecycleService) createAgent(
return fmt.Errorf("agent record creation failed: %w", err)
}
// Create default subsystems for new agent
if err := s.subsystemQueries.CreateDefaultSubsystems(agent.ID); err != nil {
// Determine platform-specific scanners based on OS type
availableScanners := s.determineScannersFromOS(cfg.Platform)
s.logger.Printf("[INFO] [server] [lifecycle] agent=%s os=%s scanners=%v", cfg.AgentID, cfg.Platform, availableScanners)
// Create platform-specific subsystems for new agent
if err := s.subsystemQueries.CreateDefaultSubsystems(agent.ID, availableScanners); err != nil {
s.logger.Printf("Warning: failed to create default subsystems: %v", err)
// Non-fatal error - agent still created
}
@ -206,6 +211,40 @@ func (s *AgentLifecycleService) createAgent(
return nil
}
// determineScannersFromOS returns appropriate package manager scanners based on OS type
func (s *AgentLifecycleService) determineScannersFromOS(osType string) []string {
scanners := []string{}
osLower := strings.ToLower(osType)
switch {
// Debian/Ubuntu based systems
case strings.Contains(osLower, "debian"), strings.Contains(osLower, "ubuntu"), strings.Contains(osLower, "linuxmint"):
scanners = append(scanners, "apt")
// Fedora/RHEL/CentOS based systems
case strings.Contains(osLower, "fedora"), strings.Contains(osLower, "rhel"), strings.Contains(osLower, "centos"), strings.Contains(osLower, "rocky"), strings.Contains(osLower, "alma"):
scanners = append(scanners, "dnf")
// Windows systems
case strings.Contains(osLower, "windows"):
scanners = append(scanners, "winget", "windows")
// Arch-based (fallback to apt for structure similarity, or could add pacman)
case strings.Contains(osLower, "arch"), strings.Contains(osLower, "manjaro"):
scanners = append(scanners, "apt") // Using apt as fallback
// Default fallback for unknown Linux
case strings.Contains(osLower, "linux"):
scanners = append(scanners, "apt") // Conservative fallback
// Unknown OS - no package scanners
default:
s.logger.Printf("[WARNING] [server] [lifecycle] unknown_os_type=%s no_package_scanners", osType)
}
return scanners
}
// updateAgent updates existing agent record
func (s *AgentLifecycleService) updateAgent(
ctx context.Context,

View file

@ -70,8 +70,8 @@ func (s *ConfigService) GenerateNewConfig(agentCfg *AgentConfig) ([]byte, error)
agentID := uuid.MustParse(agentCfg.AgentID)
subsystems, err := s.subsystemQueries.GetSubsystems(agentID)
if err != nil || len(subsystems) == 0 {
// If not found, create defaults
if err := s.subsystemQueries.CreateDefaultSubsystems(agentID); err != nil {
// If not found, create defaults (nil = generic subsystems for config path)
if err := s.subsystemQueries.CreateDefaultSubsystems(agentID, nil); err != nil {
return nil, fmt.Errorf("failed to create default subsystems: %w", err)
}
subsystems, _ = s.subsystemQueries.GetSubsystems(agentID)

View file

@ -20,7 +20,7 @@
"react": "^18.2.0",
"react-dom": "^18.2.0",
"react-hot-toast": "^2.6.0",
"react-router-dom": "^6.31.0",
"react-router-dom": "^6.30.0",
"tailwind-merge": "^2.0.0",
"zustand": "^5.0.8"
},