mirror of
https://github.com/safedep/pmg.git
synced 2026-08-03 07:24:09 +02:00
* refactor(flows): extract SetupCACertificate for reuse Move the CA load/generate/merge logic out of proxyFlow into an exported flows.SetupCACertificate so the persistent proxy server can reuse it. * feat(proxy): add persistent proxy server with start/stop/env/status Introduces 'pmg proxy' commands backed by internal/proxyserver: a long-lived MITM proxy that intercepts package managers via env vars (no shims). Supports --daemon (Unix), --state, --port; generic 'env' output that skips cert vars when the CA is OS-trusted; opt-in 'stop --fail-on-violation' (fail-closed on crash) with a synchronous cloud event flush; and the malysis analysis cache. * feat(action): add server-mode for persistent proxy When server-mode=true the action starts the proxy daemon and injects proxy env vars into the job instead of installing shims. * test(proxy): add persistent proxy server E2E workflow * docs(readme): document persistent proxy server mode * fix(proxy): create cache dir before writing state file and daemon log On a fresh CI runner the cache directory does not exist yet; os.OpenFile and os.WriteFile do not create parent dirs, so 'pmg proxy start --daemon' failed with 'no such file or directory'. MkdirAll the parent before writing. * docs: add persistent proxy server architecture doc * refactor proxyserver * fix(proxy): always emit cert env vars instead of skipping on OS-trust status npm/pip/yarn/requests trust the MITM CA inconsistently across tools, versions, and configs; many still use bundled CA stores. Always emitting the cert-path env vars is the conservative choice that works regardless, and is harmless for tools that read the OS store (they ignore the vars). Skipping them when a system CA exists would silently break any tool still on a bundled store. * refactor(proxy): drop redundant audit init in daemon; rely on main.go main.go's PersistentPreRun already initializes the audit pipeline for every command (including the daemon's re-exec'd child) and closes it at process exit. Re-initializing in proxyserver.Run created a second auditor and a second cloud-sync WAL connection, orphaning the first. Removing it makes the daemon consistent with the normal proxy flow, which never self-initializes audit. * fix(proxy): bypass proxy env when flushing events to cloud on stop pmg proxy stop inherits HTTP(S)_PROXY (injected by 'pmg proxy env') pointing at the PMG proxy it just shut down. The cloud sync gRPC client honored those vars and routed api.safedep.io through the dead proxy, failing with 'connection refused' so no events were delivered. Clear the proxy env vars before the sync so PMG's own cloud traffic goes direct. * chore(proxy): address review feedback - configurable bind host via proxy.server.listen_host (default loopback) - proxy commands use ui.ErrorExit instead of returning errors to cobra - rename errcode to ProxyPolicyViolation (covers malware + cooldown) - share cloud sync via audit.DrainToCloud (de-dup with cmd/cloud/sync) - centralize proxy CA bundle path in certmanager - docs: persistent proxy cert trust + bind address * fix(proxy): show real message on fail-on-violation error stopExitError set only WithMsg, but ui.ErrorExit renders HumanError, so the framed error showed 'no human-readable message available'. Set both from one string, and emit the framed error before the stdout summary so the blocked count is stated once. * fix(proxy): flush cloud events from the daemon, not stop The stop process inherits HTTP_PROXY (from 'pmg proxy env'), so its cloud client routed api.safedep.io through the already-stopped proxy and failed with connection refused. Move the flush into the daemon's shutdown, which has no proxy env (it started before env injection) and dials SafeDep directly. - daemon flushes on shutdown via audit.DrainToCloud and records the result in the state file; stop surfaces it (on both success and fail-on-violation paths) since the daemon's own logs aren't visible to stop - coordinate stop's wait with the daemon shutdown budget; on timeout, error out without reading stale state or deleting the file (fail-closed) - persist blocked count before the flush so the gate stays correct if the flush hangs or the daemon is killed mid-flush - remove now-redundant cloud_flush.go * disable auto-sync for proxy cmds * feat(proxy): periodic cloud sync + move proxy env vars to packagemanager - daemon runs a periodic cloud-sync ticker so the shutdown flush stays small; the run total is reported by stop, and shutdown timeouts are coordinated - move EnvVarForProxy from config to packagemanager (it is package-manager knowledge); the shared function now builds the proxy URL and NO_PROXY itself, removing the duplicated construction in the per-command and persistent paths - relocate the #319 yarn and #339 IPv6 regression tests alongside the function - enable cloud sync in the persistent-proxy E2E workflow and fix the stale internal/proxystate path filter * refactor(proxy): rename cloudFlushLockTimeout to cloudFlushLockWait Consistent timeout naming: *LockWait is the lock-acquire bound, *Timeout is the sync-RPC bound. Previously the final-flush pair was cloudFlushLockTimeout vs cloudFlushTimeout — two lookalike names for different operations. * refactor(proxy): extract cloudFlush and trim duplicate shutdown comments The shutdown's final-flush block is now a cloudFlush helper, symmetric with startCloudSyncLoop (one-shot vs loop). Removed the triplicated ticker/lock contention comments, keeping the contract on the function doc and one-line pointers at the call sites. * docs: update persistent proxy cloud sync to daemon-owned model The daemon now owns cloud delivery (periodic sync while serving + final flush on shutdown); stop signals it, waits, and reports the result. Rewrite the Cloud event sync section, fix stop attributions, add the cloud_sync state field, and update the sequence diagram. * docs: move Usage section up below How it works Put the copy-paste recipes near the top so users find them before the internals. * refactor(proxy): address PR review feedback - configurable bind host/port via --host/--port flags + config (listen_host, listen_port), bound directly to config fields per PMG's flag pattern - daemon log path via --log-file and readiness timeout in ProxyDaemonConfig; Daemonize no longer owns path policy (caller validates, fails fast) - gate periodic cloud sync on auto_sync; suppress detached background sync for proxy commands instead of flipping the flag - pmg proxy env --export emits shell-quoted lines for eval (spaces survive) - extract shared flows.BuildCachedMalysisAnalyzer, dropping the analyzer+cache duplication between proxy flow and proxy server - add internal/proxyserver/doc.go documenting the package + boundary vs flows - E2E: assert malicious installs are blocked (drop continue-on-error) - docs: trim Commands/State-file to user contracts; refresh bind address * refactor(proxy): proactive alignment fixes from whole-PR review - gate the shutdown cloud flush on auto_sync too, matching the periodic ticker (auto_sync consistently controls all daemon-driven cloud delivery) - ResolveStatePath takes cacheDir instead of *RuntimeConfig, keeping state.go free of config dependency - drop the empty-host comment in listenAddr; keep the loopback guard so a blank host never silently binds all interfaces * fix: Decouple localdb with malysis analyser construction * fix: Persist global args before proxy server daemon exec * fix: GitHub Action for cloud auto-sync in server mode --------- Co-authored-by: Abhisek Datta <abhisek.datta@gmail.com>
277 lines
9.8 KiB
Go
277 lines
9.8 KiB
Go
package proxyserver
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"fmt"
|
|
"net"
|
|
"os"
|
|
"os/signal"
|
|
"strconv"
|
|
"syscall"
|
|
"time"
|
|
|
|
"github.com/safedep/dry/log"
|
|
"github.com/safedep/pmg/config"
|
|
"github.com/safedep/pmg/internal/audit"
|
|
"github.com/safedep/pmg/internal/flows"
|
|
"github.com/safedep/pmg/internal/localstore"
|
|
pmgproxy "github.com/safedep/pmg/proxy"
|
|
"github.com/safedep/pmg/proxy/certmanager"
|
|
"github.com/safedep/pmg/proxy/interceptors"
|
|
)
|
|
|
|
const (
|
|
serverStopTimeout = 5 * time.Second
|
|
|
|
// Periodic cloud sync runs while the daemon is alive so most audit events are
|
|
// delivered during the run and the shutdown flush stays small. A tick that
|
|
// cannot get the sync lock quickly is skipped (the next tick retries).
|
|
cloudSyncInterval = 15 * time.Second
|
|
cloudSyncTickLockWait = 5 * time.Second
|
|
cloudSyncTickTimeout = 30 * time.Second
|
|
|
|
// Final flush at shutdown, when the daemon drains whatever the ticker left.
|
|
cloudFlushLockWait = 30 * time.Second
|
|
cloudFlushTimeout = 2 * time.Minute
|
|
|
|
// daemonShutdownBudget is the worst-case time the daemon needs to shut down:
|
|
// drain in-flight requests, wait for an in-flight periodic tick to finish,
|
|
// then the final flush. `pmg proxy stop` waits at least this long for the
|
|
// daemon to exit; see stopWaitTimeout.
|
|
daemonShutdownBudget = serverStopTimeout +
|
|
cloudSyncTickLockWait + cloudSyncTickTimeout +
|
|
cloudFlushLockWait + cloudFlushTimeout
|
|
)
|
|
|
|
// DefaultDaemonReadyTimeout is how long the parent waits for the daemon to
|
|
// become ready before giving up, when ProxyDaemonConfig.ReadyTimeout is unset.
|
|
const DefaultDaemonReadyTimeout = 10 * time.Second
|
|
|
|
// ProxyDaemonConfig carries the daemon-launch parameters the caller decides, so
|
|
// daemonization stays free of config and path-policy concerns.
|
|
type ProxyDaemonConfig struct {
|
|
// LogPath is the file the detached daemon's stdout/stderr is redirected to.
|
|
// The caller owns this path (its parent directory must exist).
|
|
LogPath string
|
|
// ReadyTimeout bounds how long to wait for the daemon to write its state
|
|
// file and become live.
|
|
ReadyTimeout time.Duration
|
|
}
|
|
|
|
// Run starts the persistent proxy server in the foreground and blocks until it
|
|
// receives SIGINT/SIGTERM. It writes the state file on startup, auto-blocks
|
|
// suspicious packages, and records the final blocked count on shutdown.
|
|
func Run(ctx context.Context, cfg *config.RuntimeConfig, statePath, host string, port int) error {
|
|
if existing, err := readState(statePath); err == nil && existing.IsRunning() {
|
|
return fmt.Errorf("proxy already running (pid %d, addr %s) — run 'pmg proxy stop' first", existing.PID, existing.Addr)
|
|
}
|
|
|
|
caCertPath := certmanager.ProxyCABundlePath(cfg.ConfigDir())
|
|
caCert, _, err := flows.SetupCACertificate(cfg.ConfigDir(), caCertPath)
|
|
if err != nil {
|
|
return fmt.Errorf("setup CA certificate: %w", err)
|
|
}
|
|
|
|
certMgr, err := certmanager.NewCertificateManagerWithCA(caCert, certmanager.DefaultCertManagerConfig())
|
|
if err != nil {
|
|
return fmt.Errorf("create certificate manager: %w", err)
|
|
}
|
|
|
|
localDB := localstore.NewManager(cfg)
|
|
defer func() {
|
|
if cerr := localDB.Close(); cerr != nil {
|
|
log.Warnf("failed to close localdb: %v", cerr)
|
|
}
|
|
}()
|
|
|
|
malysisAnalyzer, err := flows.BuildMalysisAnalyzer(ctx, cfg, localDB)
|
|
if err != nil {
|
|
return fmt.Errorf("create analyzer: %w", err)
|
|
}
|
|
|
|
cache := interceptors.NewInMemoryAnalysisCache()
|
|
stats := interceptors.NewAnalysisStatsCollector()
|
|
confirmationChan := make(chan *interceptors.ConfirmationRequest, 100)
|
|
go autoBlockConfirmations(confirmationChan)
|
|
|
|
factory := interceptors.NewInterceptorFactory(
|
|
malysisAnalyzer, cache, stats, confirmationChan, interceptors.InterceptorContext{},
|
|
)
|
|
|
|
var interceptorList []pmgproxy.Interceptor
|
|
for _, eco := range interceptors.SupportedEcosystems() {
|
|
i, ferr := factory.CreateInterceptor(eco)
|
|
if ferr != nil {
|
|
return fmt.Errorf("create interceptor for %s: %w", eco.String(), ferr)
|
|
}
|
|
interceptorList = append(interceptorList, i)
|
|
}
|
|
interceptorList = append(interceptorList, interceptors.NewAuditLoggerInterceptor())
|
|
|
|
proxyConfig := pmgproxy.DefaultProxyConfig()
|
|
proxyConfig.ListenAddr = listenAddr(host, port)
|
|
proxyConfig.CertManager = certMgr
|
|
proxyConfig.Interceptors = interceptorList
|
|
|
|
server, err := pmgproxy.NewProxyServer(proxyConfig)
|
|
if err != nil {
|
|
return fmt.Errorf("create proxy server: %w", err)
|
|
}
|
|
|
|
if err := server.Start(); err != nil {
|
|
return fmt.Errorf("start proxy server: %w", err)
|
|
}
|
|
|
|
state := State{
|
|
PID: os.Getpid(),
|
|
Addr: server.Address(),
|
|
CACertPath: caCertPath,
|
|
}
|
|
if err := writeState(statePath, state); err != nil {
|
|
stopCtx, cancel := context.WithTimeout(context.Background(), serverStopTimeout)
|
|
defer cancel()
|
|
if serr := server.Stop(stopCtx); serr != nil {
|
|
log.Warnf("failed to stop proxy after state write failure: %v", serr)
|
|
}
|
|
return fmt.Errorf("write proxy state: %w", err)
|
|
}
|
|
|
|
log.Infof("PMG persistent proxy running on %s (pid %d)", state.Addr, state.PID)
|
|
if _, err := fmt.Fprintf(os.Stderr, "PMG proxy running on %s\nRun: export $(pmg proxy env | xargs) # or: pmg proxy env >> \"$GITHUB_ENV\"\n", state.Addr); err != nil {
|
|
log.Warnf("failed to write startup message: %v", err)
|
|
}
|
|
|
|
// Periodically flush events to the cloud while serving so the shutdown flush
|
|
// stays small. See startCloudSyncLoop for the stop-function contract.
|
|
stopSyncLoop := startCloudSyncLoop(cfg)
|
|
|
|
sigCh := make(chan os.Signal, 1)
|
|
signal.Notify(sigCh, os.Interrupt, syscall.SIGTERM)
|
|
<-sigCh
|
|
|
|
// Drain in-flight requests before closing the confirmation channel, so no
|
|
// request handler can send on a closed channel (panic) during shutdown.
|
|
stopCtx, cancel := context.WithTimeout(context.Background(), serverStopTimeout)
|
|
defer cancel()
|
|
stopErr := server.Stop(stopCtx)
|
|
|
|
close(confirmationChan)
|
|
|
|
// Count is read after drain so a package analyzed at shutdown is not missed.
|
|
// Persist it BEFORE the (possibly slow) cloud flush so the blocked count
|
|
// survives even if the flush hangs or the daemon is killed mid-flush, which
|
|
// keeps `stop --fail-on-violation` correct in those cases.
|
|
state.BlockedCount = stats.GetStats().BlockedCount
|
|
if werr := writeState(statePath, state); werr != nil {
|
|
log.Warnf("failed to write final proxy state: %v", werr)
|
|
}
|
|
|
|
// Halt the periodic sync (waits for any in-flight drain) before the final
|
|
// flush, so the two never hold the sync lock at once.
|
|
periodicSynced := stopSyncLoop()
|
|
|
|
if cs := cloudFlush(cfg, periodicSynced); cs != nil {
|
|
state.CloudSync = cs
|
|
if werr := writeState(statePath, state); werr != nil {
|
|
log.Warnf("failed to write final proxy state: %v", werr)
|
|
}
|
|
}
|
|
|
|
return stopErr
|
|
}
|
|
|
|
// cloudFlush drains whatever the periodic sync left and returns the outcome
|
|
// (total delivered, including periodicSynced). Returns nil when automatic cloud
|
|
// delivery is off (cloud or auto-sync disabled) — same gate as the periodic
|
|
// ticker, so auto_sync consistently controls all daemon-driven cloud delivery.
|
|
// The daemon does this itself rather than `pmg proxy stop` because, unlike stop,
|
|
// it has no proxy env vars and so dials SafeDep directly instead of routing
|
|
// through the proxy that is now shutting down. Uses a fresh context since the
|
|
// caller's may already be cancelled at shutdown.
|
|
func cloudFlush(cfg *config.RuntimeConfig, periodicSynced int) *CloudSyncResult {
|
|
if !cfg.Config.Cloud.Enabled || !cfg.Config.Cloud.AutoSync.Enabled {
|
|
return nil
|
|
}
|
|
|
|
synced, err := audit.DrainToCloud(context.Background(), cfg, cloudFlushLockWait, cloudFlushTimeout)
|
|
res := &CloudSyncResult{Synced: periodicSynced + synced}
|
|
if err != nil {
|
|
res.Error = err.Error()
|
|
log.Warnf("cloud event flush failed: %v", err)
|
|
} else {
|
|
log.Infof("Flushed %d events to SafeDep Cloud (%d during the run)", res.Synced, periodicSynced)
|
|
}
|
|
return res
|
|
}
|
|
|
|
// startCloudSyncLoop periodically drains pending audit events to SafeDep Cloud
|
|
// while the daemon runs. It returns a stop function that halts the ticker, waits
|
|
// for any in-flight drain to finish, and returns the running total of events
|
|
// delivered. The stop function must be called before the daemon's final flush so
|
|
// the two never hold the sync lock at once. A no-op when cloud sync or auto-sync
|
|
// is disabled; the daemon's automatic cloud delivery (periodic ticker and
|
|
// shutdown flush alike) honors the auto_sync flag.
|
|
func startCloudSyncLoop(cfg *config.RuntimeConfig) func() int {
|
|
if !cfg.Config.Cloud.Enabled || !cfg.Config.Cloud.AutoSync.Enabled {
|
|
return func() int { return 0 }
|
|
}
|
|
|
|
stop := make(chan struct{})
|
|
done := make(chan struct{})
|
|
var total int // written only by the goroutine; read after <-done (happens-before)
|
|
|
|
go func() {
|
|
defer close(done)
|
|
ticker := time.NewTicker(cloudSyncInterval)
|
|
defer ticker.Stop()
|
|
|
|
for {
|
|
select {
|
|
case <-stop:
|
|
return
|
|
case <-ticker.C:
|
|
synced, err := audit.DrainToCloud(context.Background(), cfg, cloudSyncTickLockWait, cloudSyncTickTimeout)
|
|
if err != nil {
|
|
if errors.Is(err, audit.ErrSyncInProgress) {
|
|
log.Debugf("periodic cloud sync skipped: another sync in progress")
|
|
} else {
|
|
log.Warnf("periodic cloud sync failed: %v", err)
|
|
}
|
|
continue
|
|
}
|
|
total += synced
|
|
if synced > 0 {
|
|
log.Infof("Periodic cloud sync: flushed %d events", synced)
|
|
}
|
|
}
|
|
}
|
|
}()
|
|
|
|
return func() int {
|
|
close(stop)
|
|
<-done
|
|
return total
|
|
}
|
|
}
|
|
|
|
// listenAddr resolves the proxy's bind address from config (host) and the
|
|
// --port flag. Host defaults to loopback.
|
|
func listenAddr(host string, port int) string {
|
|
if host == "" {
|
|
host = "127.0.0.1"
|
|
}
|
|
|
|
return net.JoinHostPort(host, strconv.Itoa(port))
|
|
}
|
|
|
|
// autoBlockConfirmations drains the confirmation channel and always denies,
|
|
// appropriate for non-interactive CI/CD environments.
|
|
func autoBlockConfirmations(ch chan *interceptors.ConfirmationRequest) {
|
|
for req := range ch {
|
|
log.Warnf("Persistent proxy: auto-blocking suspicious package %s", req.PackageVersion.GetPackage().GetName())
|
|
req.ResponseChan <- false
|
|
close(req.ResponseChan)
|
|
}
|
|
}
|