Skip to content
Merged
Show file tree
Hide file tree
Changes from 6 commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 1 addition & 2 deletions agent/backend/devicediscovery/device_discovery.go
Original file line number Diff line number Diff line change
Expand Up @@ -247,8 +247,7 @@ func (d *deviceDiscoveryBackend) Start(ctx context.Context, cancelFunc context.C
}
version, readinessErr = d.Version()
if readinessErr == nil {
d.logger.Info("device-discovery readiness ok, got version ",
"device_discovery_version", version)
d.logger.Info("device-discovery readiness ok, got version", "version", version)
break
}
backoffDuration := time.Duration(backoff) * time.Second
Expand Down
3 changes: 1 addition & 2 deletions agent/backend/networkdiscovery/network_discovery.go
Original file line number Diff line number Diff line change
Expand Up @@ -264,8 +264,7 @@ func (d *networkDiscoveryBackend) Start(ctx context.Context, cancelFunc context.
}
version, readinessErr = d.Version()
if readinessErr == nil {
d.logger.Info("network-discovery readiness ok, got version ",
"network_discovery_version", version)
d.logger.Info("network-discovery readiness ok, got version", "version", version)
break
}
backoffDuration := time.Duration(backoff) * time.Second
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -181,8 +181,7 @@ func (o *openTelemetryBackend) Start(ctx context.Context, cancelFunc context.Can
}
version, readinessErr = o.Version()
if readinessErr == nil {
o.logger.Info("opentelemetry infinity readiness ok, got version ",
"opentelemetry_infinity_version", version)
o.logger.Info("opentelemetry infinity readiness ok, got version", "version", version)
break
}
backoffDuration := time.Duration(backoff) * time.Second
Expand Down
16 changes: 12 additions & 4 deletions agent/backend/pktvisor/pktvisor.go
Original file line number Diff line number Diff line change
Expand Up @@ -205,8 +205,7 @@ func (p *pktvisorBackend) Start(ctx context.Context, cancelFunc context.CancelFu
readinessError = backend.CommonRequest("pktvisor", p.proc, p.logger, url, &appMetrics, http.MethodGet,
http.NoBody, "application/json", readinessTimeout, "error")
if readinessError == nil {
p.logger.Info("pktvisor readiness ok, got version ",
"pktvisor_version", appMetrics.App.Version)
p.logger.Info("pktvisor readiness ok, got version", "version", appMetrics.App.Version)
break
}
backoffDuration := time.Duration(backoff) * time.Second
Expand Down Expand Up @@ -314,6 +313,7 @@ func parsePktvisorEntity(line string) (entity, name, rest string, ok bool) {
func (p *pktvisorBackend) Stop(ctx context.Context) error {
p.logger.Info("routine call to stop pktvisor", "routine", ctx.Value(config.ContextKey("routine")))
defer p.cancelFunc()

err := p.proc.Stop()
finalStatus := <-p.statusChan
if err != nil {
Expand All @@ -337,6 +337,15 @@ func (p *pktvisorBackend) Configure(logger *slog.Logger, repo policies.PolicyRep
p.adminAPIPort = defaultAPIPort
p.agentLabels = common.Otlp.AgentLabels

// Clean up old temp config file if it exists
if p.configFile != "" {
if err := os.Remove(p.configFile); err != nil && !os.IsNotExist(err) {
p.logger.Warn("failed to remove old pktvisor temp config file",
"file", p.configFile,
"error", err)
}
}

// Create temp config file
tmpDir := os.TempDir()
tmpFile, err := os.CreateTemp(tmpDir, "pktvisor-*.yaml")
Expand Down Expand Up @@ -431,8 +440,7 @@ func (p *pktvisorBackend) FullReset(ctx context.Context) error {
return err
}
}

// for each policy, restart the scraper
// create a new context for the backend
backendCtx, cancelFunc := context.WithCancel(context.WithValue(ctx, config.ContextKey("routine"), "pktvisor"))

// start it
Expand Down
3 changes: 1 addition & 2 deletions agent/backend/snmpdiscovery/snmp_discovery.go
Original file line number Diff line number Diff line change
Expand Up @@ -265,8 +265,7 @@ func (d *snmpDiscoveryBackend) Start(ctx context.Context, cancelFunc context.Can
}
version, readinessErr = d.Version()
if readinessErr == nil {
d.logger.Info("snmp-discovery readiness ok, got version ",
"snmp_discovery_version", version)
d.logger.Info("snmp-discovery readiness ok, got version", "version", version)
break
}
backoffDuration := time.Duration(backoff) * time.Second
Expand Down
3 changes: 1 addition & 2 deletions agent/backend/worker/worker.go
Original file line number Diff line number Diff line change
Expand Up @@ -248,8 +248,7 @@ func (d *workerBackend) Start(ctx context.Context, cancelFunc context.CancelFunc
}
version, readinessErr = d.Version()
if readinessErr == nil {
d.logger.Info("worker readiness ok, got version ",
"worker_version", version)
d.logger.Info("worker readiness ok, got version ", "version", version)
break
}
backoffDuration := time.Duration(backoff) * time.Second
Expand Down
5 changes: 2 additions & 3 deletions agent/docker/Dockerfile
Original file line number Diff line number Diff line change
Expand Up @@ -47,7 +47,7 @@ FROM python:3.14-slim-trixie

RUN \
apt update && \
apt install --yes --no-install-recommends nmap openssh-client && \
apt install --yes --no-install-recommends nmap openssh-client tini && \
rm -rf /var/lib/apt/lists/*

RUN mkdir -p /opt/orb
Expand All @@ -70,7 +70,6 @@ COPY --from=snmp-discovery /usr/local/bin/snmp-discovery /usr/local/bin/snmp-dis

COPY --from=builder /build/orb-agent /usr/local/bin/orb-agent
COPY --from=builder /src/orb-agent/agent/docker/orb-agent-entry.sh /usr/local/bin/orb-agent-entry.sh
COPY --from=builder /src/orb-agent/agent/docker/run-agent.sh /run-agent.sh
COPY --from=builder /src/orb-agent/agent/docker/default_config.yaml /opt/orb/default_config.yaml

ENTRYPOINT [ "/usr/local/bin/orb-agent-entry.sh" ]
ENTRYPOINT [ "/usr/bin/tini", "--", "/usr/local/bin/orb-agent-entry.sh" ]
40 changes: 3 additions & 37 deletions agent/docker/orb-agent-entry.sh
Original file line number Diff line number Diff line change
Expand Up @@ -3,18 +3,6 @@
# entry point for orb-agent
#

agentstop1 () {
printf "\rFinishing container.."
exit 0
}

agentstop2 () {
if [ -f "/var/run/orb-agent.pid" ]; then
ID=$(cat /var/run/orb-agent.pid)
kill -15 $ID
fi
}

if [ "${INSTALL_DRIVERS_PATH}" != '' ]; then
cd "$(dirname "$(realpath "${INSTALL_DRIVERS_PATH}")")"
echo "Installing additional drivers"
Expand Down Expand Up @@ -74,28 +62,6 @@ if [ -n "${FLEET_CLIENT_ID}" ] && [ -n "${FLEET_CLIENT_SECRET}" ]; then
fi
fi

trap agentstop1 SIGINT
trap agentstop2 SIGTERM

# eternal loop
while true
do
# pid file dont exist
if [ ! -f "/var/run/orb-agent.pid" ]; then
# running orb-agent in background
nohup /run-agent.sh "${agent_args[@]}" &
sleep 2
if [ -d "/nohup.out" ]; then
tail -f /nohup.out &
fi
else
PID=$(cat /var/run/orb-agent.pid)
if [ ! -d "/proc/$PID" ]; then
# stop container
echo "$PID is not running"
rm /var/run/orb-agent.pid
exit 1
fi
sleep 5
fi
done
# Use exec to replace this shell process with the agent
# This makes the agent a direct child of tini, ensuring proper signal handling
exec /usr/local/bin/orb-agent "${agent_args[@]}"
Comment thread
leoparente marked this conversation as resolved.
12 changes: 0 additions & 12 deletions agent/docker/run-agent.sh

This file was deleted.

5 changes: 3 additions & 2 deletions agent/otlpbridge/server.go
Original file line number Diff line number Diff line change
Expand Up @@ -104,8 +104,9 @@ func (s *BridgeServer) GetPolicyRepo() policies.PolicyRepo {

// Start starts the gRPC server without establishing MQTT.
// Publisher and topic should be set before OTLP data arrives.
func (s *BridgeServer) Start(_ context.Context) error {
lis, err := net.Listen("tcp", s.cfg.ListenAddr)
func (s *BridgeServer) Start(ctx context.Context) error {
// Platform-specific socket configuration (SO_REUSEADDR on Unix for faster port reuse)
lis, err := listen(ctx, s.cfg.ListenAddr)
if err != nil {
return fmt.Errorf("failed to listen on %s (port may be in use by another service): %w", s.cfg.ListenAddr, err)
}
Expand Down
34 changes: 34 additions & 0 deletions agent/otlpbridge/socket_unix.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,34 @@
//go:build unix

package otlpbridge

import (
"context"
"net"
"syscall"

"golang.org/x/sys/unix"
Comment thread
leoparente marked this conversation as resolved.
)

// newListenConfig returns a net.ListenConfig with SO_REUSEADDR enabled for faster port reuse.
// This is particularly important for docker restart scenarios where ports may be in TIME_WAIT.
func newListenConfig() net.ListenConfig {
return net.ListenConfig{
Control: func(_, _ string, c syscall.RawConn) error {
var sockOptErr error
if err := c.Control(func(fd uintptr) {
// Enable SO_REUSEADDR to allow binding to TIME_WAIT sockets
sockOptErr = unix.SetsockoptInt(int(fd), unix.SOL_SOCKET, unix.SO_REUSEADDR, 1)
}); err != nil {
return err
}
return sockOptErr
},
}
}

// listen creates a TCP listener with platform-specific socket options.
func listen(ctx context.Context, addr string) (net.Listener, error) {
lc := newListenConfig()
return lc.Listen(ctx, "tcp", addr)
}
17 changes: 17 additions & 0 deletions agent/otlpbridge/socket_windows.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,17 @@
//go:build windows

package otlpbridge

import (
"context"
"net"
)

// listen creates a TCP listener using standard net.Listen.
// Windows doesn't need SO_REUSEADDR configuration like Unix systems do.
func listen(ctx context.Context, addr string) (net.Listener, error) {
// On Windows, we use standard Listen without SO_REUSEADDR
// Windows handles port reuse differently and doesn't have the same TIME_WAIT issues
var lc net.ListenConfig
return lc.Listen(ctx, "tcp", addr)
}
1 change: 1 addition & 0 deletions cmd/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -113,6 +113,7 @@ func Run(_ *cobra.Command, _ []string) {
logger.Warn("stop signal received stopping agent")
a.Stop(rootCtx)
cancelFunc()
os.Exit(0) // Exit after clean shutdown
case <-rootCtx.Done():
logger.Warn("mainRoutine context cancelled")
done <- true
Expand Down
Loading