diff --git a/agent.yml b/agent.yml index 231ad73..5f744e6 100644 --- a/agent.yml +++ b/agent.yml @@ -1,8 +1,9 @@ env: WEB_BINDING: "0.0.0.0:23000" MANAGED: false - REMOTE_CONFIG_SERVERS: ["http://localhost:22000"] + REMOTE_CONFIG_SERVERS: ["http://localhost:29000"] REMOTE_CONFIG_INTERVAL: "1s" + ENROLLMENT_TOKEN: "" SECURITY_ENABLED: true SECURITY_MANAGED_ENABLED: false @@ -11,6 +12,10 @@ path.logs: log path.configs: "config" configs.auto_reload: true + +transfer: + receiver_port_range: "30000-30100" + resource_limit.cpu.max_num_of_cpus: 1 resource_limit: memory: @@ -49,7 +54,6 @@ api: enabled: false web: - embedding_api: false enabled: true network: binding: $[[env.WEB_BINDING]] @@ -67,8 +71,6 @@ web: access_token: enabled: true -agent: - metrics: enabled: true @@ -79,6 +81,7 @@ configs: allow_generated_metrics_tasks: false # allow auto-generated metrics tasks (e.g. k8s) interval: $[[env.REMOTE_CONFIG_INTERVAL]] servers: $[[env.REMOTE_CONFIG_SERVERS]] # config servers + enrollment_token: $[[env.ENROLLMENT_TOKEN]] # one-time registration pass (required when the server sets enrollment.required; pass via -e ENROLLMENT_TOKEN=et-...) max_backup_files: 5 soft_delete: false # tls: #for mTLS connection with config servers @@ -86,3 +89,142 @@ configs: # cert_file: /etc/ssl.crt # key_file: /etc/ssl.key # skip_insecure_verify: false + +# --------------------------------------------------------------------------- +# Elasticsearch/Easysearch log viewing whitelist (optional). +# +# /elasticsearch/logs/_list and /elasticsearch/logs/_read only serve +# directories from this whitelist plus the log paths the local search +# nodes report themselves (settings path.logs). System paths (/etc, +# /usr/bin, ...) are never readable. Only needed when discovery cannot +# see the log directory (e.g. logs mounted elsewhere); changes apply +# within a minute or after restart. +# --------------------------------------------------------------------------- +#elasticsearch_logs: +# allowed_paths: +# - /var/log/easysearch + +# --------------------------------------------------------------------------- +# Log pipeline egress (opt-in): process collected logs and ship them to +# the gateway's OTLP/gRPC intake. +# +# Flow: logs_processor writes envelopes to the local "logs" disk queue -> +# consumer drains batches -> for_each splits records and runs the +# transform chain (dissect + field_standardize) in place -> +# otlp_export ships one OTLP ExportLogsServiceRequest to the gateway. +# +# Enable by setting enabled:true and pointing endpoint at the gateway. +# --------------------------------------------------------------------------- +pipeline: + - name: logs_otlp_egress + enabled: false + auto_start: true + keep_running: true + singleton: true + processor: + - consumer: + queue_selector: + keys: + - logs + consumer: + group: agent-logs-otlp + message_field: messages + processor: + - for_each: + message_field: messages + processor: + - dissect: + pattern: "%{log_level} %{service_name}: %{remainder}" + ignore_failure: true + - field_standardize: + mode: underscore + normalize_case: true + - otlp_export: + endpoint: "127.0.0.1:4317" + insecure: true + timeout: 10s + message_field: messages + +# --------------------------------------------------------------------------- +# Syslog collection (opt-in): listen for RFC3164/RFC5424 messages and +# feed them onto the local "logs" queue in the same envelope as file +# collection, so the shared processing chain (consumer -> for_each -> +# dissect/field_standardize -> otlp_export) applies unchanged. +# +# Point remote syslog senders at this address, e.g. in rsyslog: +# *.* @127.0.0.1:5140 # UDP +# *.* @@127.0.0.1:5140 # TCP +# --------------------------------------------------------------------------- + - name: syslog_intake + enabled: false + auto_start: true + keep_running: true + processor: + - syslog_processor: + protocol: udp + bind: ":5140" + queue_name: logs + +# --------------------------------------------------------------------------- +# Kafka output (opt-in): route the local queues to Kafka instead of the +# disk queue. With kafka_queue.default=true the existing queue.Push +# calls publish to Kafka unchanged — no code change needed. +# +# disk_queue: +# default: false +# kafka_queue: +# enabled: true +# default: true +# brokers: +# - "127.0.0.1:9092" +# num_of_partition: 1 +# num_of_replica: 1 +# --------------------------------------------------------------------------- + +# --------------------------------------------------------------------------- +# File collection follow-tail (opt-in): set scan_interval on the +# logs_processor to keep re-scanning the logs path while the agent runs +# (close-to-live tailing plus rename/move-safe offset inheritance): +# +# - logs_processor: +# scan_interval: 2s +# patterns: +# - pattern: '\.log$' +# type: text +# --------------------------------------------------------------------------- + +# --------------------------------------------------------------------------- +# Direct-ship mode (opt-in): file sources can bypass the local queue +# entirely — the file itself plus offset checkpoints provide durability +# (offsets only advance after successful delivery), so the same bytes +# are not written to disk twice. Envelopes are batched in memory +# (ship_batch_size / ship_flush_interval) and shipped straight to the +# gateway's OTLP intake with load balancing, failover and retries. +# ship_config takes the same keys as the otlp_export processor. +# +# The consumer -> for_each -> otlp_export egress pipeline above is NOT +# needed for a ship_direct source; edge processors can still run via +# the queue path by leaving ship_direct off. +# +# pipeline: +# - name: logs_intake +# auto_start: true +# keep_running: true +# processor: +# - logs: +# queue_name: logs +# ship_direct: true +# ship_batch_size: 500 # events per in-flight batch +# ship_flush_interval: 1s # flush at least this often +# ship_config: +# endpoints: +# - "127.0.0.1:4317" +# insecure: true +# timeout: 10s +# scan_interval: 2s +# logs_path: /var/log/myapp +# patterns: +# - pattern: '\.log$' +# type: text +# --------------------------------------------------------------------------- + diff --git a/go.mod b/go.mod index 4e0b8a8..847b00e 100644 --- a/go.mod +++ b/go.mod @@ -7,16 +7,11 @@ replace infini.sh/framework => ../framework replace github.com/cihub/seelog => ../framework/lib/seelog require ( - github.com/cespare/xxhash/v2 v2.3.0 github.com/cihub/seelog v0.0.0-00010101000000-000000000000 - github.com/dlclark/regexp2 v1.12.0 github.com/pkg/errors v0.9.1 github.com/shirou/gopsutil/v4 v4.26.3 github.com/stretchr/testify v1.11.1 - golang.org/x/crypto v0.53.0 - golang.org/x/sys v0.47.0 golang.org/x/text v0.38.0 - gopkg.in/yaml.v3 v3.0.1 infini.sh/framework v0.0.0-00010101000000-000000000000 ) @@ -27,15 +22,18 @@ require ( github.com/arl/statsviz v0.6.0 // indirect github.com/bits-and-blooms/bitset v1.12.0 // indirect github.com/bkaradzic/go-lz4 v1.0.0 // indirect - github.com/buger/jsonparser v1.1.2 // indirect + github.com/buger/jsonparser v1.2.0 // indirect github.com/caddyserver/certmagic v0.25.3 // indirect github.com/caddyserver/zerossl v0.1.5 // indirect + github.com/cespare/xxhash/v2 v2.3.0 // indirect github.com/davecgh/go-spew v1.1.2-0.20180830191138-d8f796af33cc // indirect github.com/dgraph-io/ristretto v0.2.0 // indirect + github.com/dlclark/regexp2 v1.12.0 // indirect github.com/dustin/go-humanize v1.0.1 // indirect github.com/ebitengine/purego v0.10.0 // indirect github.com/emirpasic/gods v1.18.1 // indirect - github.com/fsnotify/fsnotify v1.9.0 // indirect + github.com/fsnotify/fsnotify v1.10.1 // indirect + github.com/go-jose/go-jose/v4 v4.1.4 // indirect github.com/go-ole/go-ole v1.2.6 // indirect github.com/golang-jwt/jwt/v4 v4.5.2 // indirect github.com/golang/protobuf v1.5.4 // indirect @@ -56,16 +54,17 @@ require ( github.com/josharian/intern v1.0.0 // indirect github.com/kardianos/osext v0.0.0-20190222173326-2bc1f35cddc0 // indirect github.com/kardianos/service v1.2.2 // indirect - github.com/klauspost/compress v1.18.0 // indirect + github.com/klauspost/compress v1.18.7 // indirect github.com/klauspost/cpuid/v2 v2.3.0 // indirect github.com/libdns/libdns v1.1.1 // indirect github.com/lufia/plan9stats v0.0.0-20211012122336-39d0f177ccd0 // indirect - github.com/mailru/easyjson v0.9.0 // indirect + github.com/mailru/easyjson v0.9.2 // indirect github.com/mark3labs/mcp-go v0.57.0 // indirect github.com/mholt/acmez/v3 v3.1.6 // indirect github.com/miekg/dns v1.1.72 // indirect github.com/mitchellh/mapstructure v1.5.0 // indirect github.com/mschoch/smat v0.2.0 // indirect + github.com/pierrec/lz4/v4 v4.1.26 // indirect github.com/pmezard/go-difflib v1.0.1-0.20181226105442-5d4384ee4fb2 // indirect github.com/power-devops/perfstat v0.0.0-20240221224432-82ca36839d55 // indirect github.com/r3labs/diff/v2 v2.15.1 // indirect @@ -73,12 +72,15 @@ require ( github.com/rs/xid v1.6.0 // indirect github.com/ryanuber/go-glob v1.0.0 // indirect github.com/santhosh-tekuri/jsonschema/v6 v6.0.2 // indirect - github.com/savsgio/gotils v0.0.0-20250408102913-196191ec6287 // indirect + github.com/savsgio/gotils v0.0.0-20250924091648-bce9a52d7761 // indirect github.com/segmentio/asm v1.1.3 // indirect github.com/segmentio/encoding v0.4.1 // indirect github.com/spf13/cast v1.7.1 // indirect github.com/tklauser/go-sysconf v0.3.16 // indirect github.com/tklauser/numcpus v0.11.0 // indirect + github.com/twmb/franz-go v1.21.6 // indirect + github.com/twmb/franz-go/pkg/kadm v1.18.0 // indirect + github.com/twmb/franz-go/pkg/kmsg v1.13.1 // indirect github.com/valyala/bytebufferpool v1.0.0 // indirect github.com/vmihailenco/msgpack v4.0.4+incompatible // indirect github.com/yosida95/uritemplate/v3 v3.0.2 // indirect @@ -87,13 +89,16 @@ require ( go.uber.org/multierr v1.11.0 // indirect go.uber.org/zap v1.27.1 // indirect go.uber.org/zap/exp v0.3.0 // indirect + golang.org/x/crypto v0.53.0 // indirect golang.org/x/mod v0.37.0 // indirect golang.org/x/net v0.56.0 // indirect golang.org/x/oauth2 v0.36.0 // indirect golang.org/x/sync v0.21.0 // indirect + golang.org/x/sys v0.47.0 // indirect golang.org/x/term v0.44.0 // indirect golang.org/x/time v0.11.0 // indirect golang.org/x/tools v0.47.0 // indirect google.golang.org/appengine v1.6.6 // indirect google.golang.org/protobuf v1.36.11 // indirect + gopkg.in/yaml.v3 v3.0.1 // indirect ) diff --git a/go.sum b/go.sum index a2025a6..d5617eb 100644 --- a/go.sum +++ b/go.sum @@ -12,8 +12,8 @@ github.com/bits-and-blooms/bitset v1.12.0 h1:U/q1fAF7xXRhFCrhROzIfffYnu+dlS38vCZ github.com/bits-and-blooms/bitset v1.12.0/go.mod h1:7hO7Gc7Pp1vODcmWvKMRA9BNmbv6a/7QIWpPxHddWR8= github.com/bkaradzic/go-lz4 v1.0.0 h1:RXc4wYsyz985CkXXeX04y4VnZFGG8Rd43pRaHsOXAKk= github.com/bkaradzic/go-lz4 v1.0.0/go.mod h1:0YdlkowM3VswSROI7qDxhRvJ3sLhlFrRRwjwegp5jy4= -github.com/buger/jsonparser v1.1.2 h1:frqHqw7otoVbk5M8LlE/L7HTnIq2v9RX6EJ48i9AxJk= -github.com/buger/jsonparser v1.1.2/go.mod h1:6RYKKt7H4d4+iWqouImQ9R2FZql3VbhNgx27UK13J/0= +github.com/buger/jsonparser v1.2.0 h1:4EFcvK1kD4jyj6YqNK6skK6w+y7FHHBR+XBCtxwu/6g= +github.com/buger/jsonparser v1.2.0/go.mod h1:6RYKKt7H4d4+iWqouImQ9R2FZql3VbhNgx27UK13J/0= github.com/caddyserver/certmagic v0.25.3 h1:mGf5ba8F7xA4c5jfDZZbK2buY1VEkbnwpMDixaju94A= github.com/caddyserver/certmagic v0.25.3/go.mod h1:YVs43D5+H/Dckt4bTga1KSO/xYfFBfVZainGDywYPAA= github.com/caddyserver/zerossl v0.1.5 h1:dkvOjBAEEtY6LIGAHei7sw2UgqSD6TrWweXpV7lvEvE= @@ -37,10 +37,10 @@ github.com/emirpasic/gods v1.18.1 h1:FXtiHYKDGKCW2KzwZKx0iC0PQmdlorYgdFG9jPXJ1Bc github.com/emirpasic/gods v1.18.1/go.mod h1:8tpGGwCnJ5H4r6BWwaV6OrWmMoPhUl5jm/FMNAnJvWQ= github.com/frankban/quicktest v1.14.6 h1:7Xjx+VpznH+oBnejlPUj8oUpdxnVs4f8XU8WnHkI4W8= github.com/frankban/quicktest v1.14.6/go.mod h1:4ptaffx2x8+WTWXmUCuVU6aPUX1/Mz7zb5vbUoiM6w0= -github.com/fsnotify/fsnotify v1.9.0 h1:2Ml+OJNzbYCTzsxtv8vKSFD9PbJjmhYF14k/jKC7S9k= -github.com/fsnotify/fsnotify v1.9.0/go.mod h1:8jBTzvmWwFyi3Pb8djgCCO5IBqzKJ/Jwo8TRcHyHii0= -github.com/go-jose/go-jose/v4 v4.1.3 h1:CVLmWDhDVRa6Mi/IgCgaopNosCaHz7zrMeF9MlZRkrs= -github.com/go-jose/go-jose/v4 v4.1.3/go.mod h1:x4oUasVrzR7071A4TnHLGSPpNOm2a21K9Kf04k1rs08= +github.com/fsnotify/fsnotify v1.10.1 h1:b0/UzAf9yR5rhf3RPm9gf3ehBPpf0oZKIjtpKrx59Ho= +github.com/fsnotify/fsnotify v1.10.1/go.mod h1:TLheqan6HD6GBK6PrDWyDPBaEV8LspOxvPSjC+bVfgo= +github.com/go-jose/go-jose/v4 v4.1.4 h1:moDMcTHmvE6Groj34emNPLs/qtYXRVcd6S7NHbHz3kA= +github.com/go-jose/go-jose/v4 v4.1.4/go.mod h1:x4oUasVrzR7071A4TnHLGSPpNOm2a21K9Kf04k1rs08= github.com/go-ole/go-ole v1.2.6 h1:/Fpf6oFPoeFik9ty7siob0G6Ke8QvQEuVcuChpwXzpY= github.com/go-ole/go-ole v1.2.6/go.mod h1:pprOEPIfldk/42T2oK7lQ4v4JSDwmV0As9GaiUsvbm0= github.com/go-viper/mapstructure/v2 v2.2.1 h1:ZAaOCxANMuZx5RCeg0mBdEZk7DZasvvZIxtHqx8aGss= @@ -88,8 +88,8 @@ github.com/kardianos/osext v0.0.0-20190222173326-2bc1f35cddc0 h1:iQTw/8FWTuc7uia github.com/kardianos/osext v0.0.0-20190222173326-2bc1f35cddc0/go.mod h1:1NbS8ALrpOvjt0rHPNLyCIeMtbizbir8U//inJ+zuB8= github.com/kardianos/service v1.2.2 h1:ZvePhAHfvo0A7Mftk/tEzqEZ7Q4lgnR8sGz4xu1YX60= github.com/kardianos/service v1.2.2/go.mod h1:CIMRFEJVL+0DS1a3Nx06NaMn4Dz63Ng6O7dl0qH0zVM= -github.com/klauspost/compress v1.18.0 h1:c/Cqfb0r+Yi+JtIEq73FWXVkRonBlf0CRNYc8Zttxdo= -github.com/klauspost/compress v1.18.0/go.mod h1:2Pp+KzxcywXVXMr50+X0Q/Lsb43OQHYWRCY2AiWywWQ= +github.com/klauspost/compress v1.18.7 h1:aUyZsS4kH3QTKurYhAOwAHxllVPnOthb3vPfnF1Ehjw= +github.com/klauspost/compress v1.18.7/go.mod h1:cwPg85FWrGar70rWktvGQj8/hthj3wpl0PGDogxkrSQ= github.com/klauspost/cpuid/v2 v2.3.0 h1:S4CRMLnYUhGeDFDqkGriYKdfoFlDnMtqTiI/sFzhA9Y= github.com/klauspost/cpuid/v2 v2.3.0/go.mod h1:hqwkgyIinND0mEev00jJYCxPNVRVXFQeu1XKlok6oO0= github.com/kr/pretty v0.3.1 h1:flRD4NNwYAUpkphVc1HcthR4KEIFJ65n8Mw5qdRn3LE= @@ -106,8 +106,8 @@ github.com/lufia/plan9stats v0.0.0-20211012122336-39d0f177ccd0 h1:6E+4a0GO5zZEnZ github.com/lufia/plan9stats v0.0.0-20211012122336-39d0f177ccd0/go.mod h1:zJYVVT2jmtg6P3p1VtQj7WsuWi/y4VnjVBn7F8KPB3I= github.com/magiconair/properties v1.8.10 h1:s31yESBquKXCV9a/ScB3ESkOjUYYv+X0rg8SYxI99mE= github.com/magiconair/properties v1.8.10/go.mod h1:Dhd985XPs7jluiymwWYZ0G4Z61jb3vdS329zhj2hYo0= -github.com/mailru/easyjson v0.9.0 h1:PrnmzHw7262yW8sTBwxi1PdJA3Iw/EKBa8psRf7d9a4= -github.com/mailru/easyjson v0.9.0/go.mod h1:1+xMtQp2MRNVL/V1bOzuP3aP8VNwRW55fQUto+XFtTU= +github.com/mailru/easyjson v0.9.2 h1:dX8U45hQsZpxd80nLvDGihsQ/OxlvTkVUXH2r/8cb2M= +github.com/mailru/easyjson v0.9.2/go.mod h1:1+xMtQp2MRNVL/V1bOzuP3aP8VNwRW55fQUto+XFtTU= github.com/mark3labs/mcp-go v0.57.0 h1:jzWKyCzdWnwnZt05cvcQQ+ngiUl2RnixXJa7Kj4qP1E= github.com/mark3labs/mcp-go v0.57.0/go.mod h1:+8WclSK1ZUweCP3hvktSji8n8ABG/95QaEkeVE/Uwas= github.com/mattn/go-isatty v0.0.24 h1:tGZZoVgT/KiqK1c8ocVLeDS8BSWMRd47J3Lbz7vsReI= @@ -126,6 +126,8 @@ github.com/nsqio/nsq v1.3.0 h1:v7NtyO844ieTIOCQEqQ7IUSSi1ImhgrTTto1rgIYGEU= github.com/nsqio/nsq v1.3.0/go.mod h1:RxNr6UC0kSkNF44LnJrlN3U3CQnQGTXk+QKfSZLzqvc= github.com/pelletier/go-toml/v2 v2.2.3 h1:YmeHyLY8mFWbdkNWwpr+qIL2bEqT0o95WSdkNHvL12M= github.com/pelletier/go-toml/v2 v2.2.3/go.mod h1:MfCQTFTvCcUyyvvwm1+G6H/jORL20Xlb6rzQu9GuUkc= +github.com/pierrec/lz4/v4 v4.1.26 h1:GrpZw1gZttORinvzBdXPUXATeqlJjqUG/D87TKMnhjY= +github.com/pierrec/lz4/v4 v4.1.26/go.mod h1:EoQMVJgeeEOMsCqCzqFm2O0cJvljX2nGZjcRIPL34O4= github.com/pkg/errors v0.9.1 h1:FEBLx1zS214owpjy7qsBeixbURkuhQAwrK5UwLGTwt4= github.com/pkg/errors v0.9.1/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0= github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= @@ -149,8 +151,8 @@ github.com/sagikazarmark/locafero v0.7.0 h1:5MqpDsTGNDhY8sGp0Aowyf0qKsPrhewaLSsF github.com/sagikazarmark/locafero v0.7.0/go.mod h1:2za3Cg5rMaTMoG/2Ulr9AwtFaIppKXTRYnozin4aB5k= github.com/santhosh-tekuri/jsonschema/v6 v6.0.2 h1:KRzFb2m7YtdldCEkzs6KqmJw4nqEVZGK7IN2kJkjTuQ= github.com/santhosh-tekuri/jsonschema/v6 v6.0.2/go.mod h1:JXeL+ps8p7/KNMjDQk3TCwPpBy0wYklyWTfbkIzdIFU= -github.com/savsgio/gotils v0.0.0-20250408102913-196191ec6287 h1:qIQ0tWF9vxGtkJa24bR+2i53WBCz1nW/Pc47oVYauC4= -github.com/savsgio/gotils v0.0.0-20250408102913-196191ec6287/go.mod h1:sM7Mt7uEoCeFSCBM+qBrqvEo+/9vdmj19wzp3yzUhmg= +github.com/savsgio/gotils v0.0.0-20250924091648-bce9a52d7761 h1:McifyVxygw1d67y6vxUqls2D46J8W9nrki9c8c0eVvE= +github.com/savsgio/gotils v0.0.0-20250924091648-bce9a52d7761/go.mod h1:Vi9gvHvTw4yCUHIznFl5TPULS7aXwgaTByGeBY75Wko= github.com/segmentio/asm v1.1.3 h1:WM03sfUOENvvKexOLp+pCqgb/WDjsi7EK8gIsICtzhc= github.com/segmentio/asm v1.1.3/go.mod h1:Ld3L4ZXGNcSLRg4JBsZ3//1+f/TjYl0Mzen/DQy1EJg= github.com/segmentio/encoding v0.4.1 h1:KLGaLSW0jrmhB58Nn4+98spfvPvmo4Ci1P/WIQ9wn7w= @@ -180,6 +182,12 @@ github.com/tklauser/go-sysconf v0.3.16 h1:frioLaCQSsF5Cy1jgRBrzr6t502KIIwQ0MArYI github.com/tklauser/go-sysconf v0.3.16/go.mod h1:/qNL9xxDhc7tx3HSRsLWNnuzbVfh3e7gh/BmM179nYI= github.com/tklauser/numcpus v0.11.0 h1:nSTwhKH5e1dMNsCdVBukSZrURJRoHbSEQjdEbY+9RXw= github.com/tklauser/numcpus v0.11.0/go.mod h1:z+LwcLq54uWZTX0u/bGobaV34u6V7KNlTZejzM6/3MQ= +github.com/twmb/franz-go v1.21.6 h1:+v0dQJVIIuw9uPmPWmPrkoUHs1pPeV8MSwA4eU/Y2kY= +github.com/twmb/franz-go v1.21.6/go.mod h1:wMepkgCatAdV9vCsuwM+wr+C1fl7KV/41+uHGAjt/wc= +github.com/twmb/franz-go/pkg/kadm v1.18.0 h1:WRf/LZmDdcDXwX7WMbtDU++v+b3NzYh2bCGoPMmzirw= +github.com/twmb/franz-go/pkg/kadm v1.18.0/go.mod h1:XeLhGoLXLFzK8/ryv5FfpxPxGwj4oFEGpPJMB/x6KDE= +github.com/twmb/franz-go/pkg/kmsg v1.13.1 h1:fG5kItwysTk5UXqVwb64EpQEy3TydF3vYYK21nUQ+bI= +github.com/twmb/franz-go/pkg/kmsg v1.13.1/go.mod h1:+DPt4NC8RmI6hqb8G09+3giKObE6uD2Eya6CfqBpeJY= github.com/valyala/bytebufferpool v1.0.0 h1:GqA5TC/0021Y/b9FG4Oi9Mr3q7XYx6KllzawFIhcdPw= github.com/valyala/bytebufferpool v1.0.0/go.mod h1:6bBcMArwyJ5K/AmCkWv1jt77kVWyCJ6HpOuEn7z0Csc= github.com/vmihailenco/msgpack v4.0.4+incompatible h1:dSLoQfGFAo3F6OoNhwUmLwVgaUXK79GlxNBwueZn0xI= diff --git a/internal/managed/exchange.go b/internal/managed/exchange.go new file mode 100644 index 0000000..2799d85 --- /dev/null +++ b/internal/managed/exchange.go @@ -0,0 +1,151 @@ +package managed + +import ( + "fmt" + "strings" + + log "github.com/cihub/seelog" + "infini.sh/framework/core/global" + "infini.sh/framework/core/keystore" + "infini.sh/framework/core/kv" + "infini.sh/framework/core/model" + "infini.sh/framework/core/security" + "infini.sh/framework/core/util" + ucfg "infini.sh/framework/lib/go-ucfg" + "infini.sh/framework/modules/configs/client" + "infini.sh/framework/modules/security/access_token" +) + +const ( + tokenExchangeAPI = "/instance/_exchange_token" + agentAPIAccessTokenKey = "AGENT_API_ACCESS_TOKEN" + managerAccessTokenKey = "CONFIGS_MANAGER_ACCESS_TOKEN" +) + +type tokenExchangeRequest struct { + InstanceID string `json:"instance_id,omitempty"` + AgentAPIToken string `json:"agent_api_token,omitempty"` +} + +type tokenExchangeResponse struct { + ManagerAPIToken string `json:"manager_api_token,omitempty"` +} + +func getOrCreateAgentAPIToken(instance model.Instance) (string, error) { + required := security.InstanceOpsPermissionKeys() + if tokenBytes, err := keystore.GetValue(agentAPIAccessTokenKey); err == nil { + token := strings.TrimSpace(string(tokenBytes)) + // Reuse the stored token only while it still carries every + // permission the web-port operational routes require; older + // tokens were minted with none, so re-mint and let the exchange + // flow register the fresh token with the manager. + if token != "" && agentAPITokenHasPermissions(token, required) { + return token, nil + } + } + + user := &security.UserSessionInfo{ + Provider: "managed_agent", + Login: instance.ID, + } + user.SetUserID(instance.ID) + user.Set("instance_id", instance.ID) + user.Set("instance_name", instance.Name) + user.Set("endpoint", instance.Endpoint) + + res, err := access_token.CreateAPIToken(user, fmt.Sprintf("%s agent api", instance.Name), "agent api token", "managed_agent_api", -1, required) + if err != nil { + return "", err + } + token, ok := res["access_token"].(string) + if !ok || strings.TrimSpace(token) == "" { + return "", fmt.Errorf("failed to create agent api access token") + } + if err := keystore.SetValue(agentAPIAccessTokenKey, []byte(token)); err != nil { + return "", err + } + return token, nil +} + +// agentAPITokenHasPermissions reports whether the stored token record still +// grants every required permission key. +func agentAPITokenHasPermissions(token string, required []security.PermissionKey) bool { + if token == "" || len(required) == 0 { + return false + } + bytes, err := kv.GetValue(access_token.KVAccessTokenBucket, []byte(token)) + if err != nil { + return false + } + record := security.AccessToken{} + if err := util.FromJSONBytes(bytes, &record); err != nil { + return false + } + granted := make(map[security.PermissionKey]bool, len(record.Permissions)) + for _, key := range record.Permissions { + granted[key] = true + } + for _, key := range required { + if !granted[key] { + return false + } + } + return true +} + +func ExchangeTokens() error { + if !global.Env().SystemConfig.Configs.Managed || len(global.Env().SystemConfig.Configs.Servers) == 0 { + return nil + } + managerAPIToken := strings.TrimSpace(global.Env().SystemConfig.Configs.ManagerConfig.AccessToken.Get()) + if managerAPIToken == "" { + return nil + } + + instance := model.GetInstanceInfo() + agentAPIToken, err := getOrCreateAgentAPIToken(instance) + if err != nil { + return err + } + + reqBody := tokenExchangeRequest{ + InstanceID: instance.ID, + AgentAPIToken: agentAPIToken, + } + req := util.Request{ + Method: util.Verb_POST, + Path: tokenExchangeAPI, + ContentType: "application/json", + Body: util.MustToJSONBytes(reqBody), + } + server, res, err := client.DoManagerRequest(&req) + if err != nil { + return err + } + if res == nil { + return fmt.Errorf("empty response from %s", server) + } + if res.StatusCode != 200 { + return fmt.Errorf("token exchange failed on %s, status: %d, body: %s", server, res.StatusCode, string(res.Body)) + } + + resp := tokenExchangeResponse{} + if err := util.FromJSONBytes(res.Body, &resp); err != nil { + return err + } + if strings.TrimSpace(resp.ManagerAPIToken) == "" { + return fmt.Errorf("manager api token is empty") + } + if err := keystore.SetValue(managerAccessTokenKey, []byte(resp.ManagerAPIToken)); err != nil { + return err + } + global.Env().SystemConfig.Configs.ManagerConfig.AccessToken = ucfg.SecretString(resp.ManagerAPIToken) + log.Infof("exchanged managed access token from %s", server) + return nil +} + +func RegisterTokenExchangeCallback() { + client.AddPostRegisterHook(func(_ string, _ *util.Result) error { + return ExchangeTokens() + }) +} diff --git a/internal/managed/module.go b/internal/managed/module.go new file mode 100644 index 0000000..911e746 --- /dev/null +++ b/internal/managed/module.go @@ -0,0 +1,27 @@ +package managed + +import "infini.sh/framework/core/module" + +const moduleName = "managed_agent_bootstrap" + +func init() { + module.RegisterModuleWithPriority(&ManagedModule{}, 100) +} + +type ManagedModule struct{} + +func (m *ManagedModule) Name() string { + return moduleName +} + +func (m *ManagedModule) Setup() { + RegisterTokenExchangeCallback() +} + +func (m *ManagedModule) Start() error { + return nil +} + +func (m *ManagedModule) Stop() error { + return nil +} diff --git a/lib/reader/harvester/harvester.go b/lib/reader/harvester/harvester.go index 29034b0..b0368f1 100644 --- a/lib/reader/harvester/harvester.go +++ b/lib/reader/harvester/harvester.go @@ -5,9 +5,11 @@ package harvester import ( + "compress/gzip" "fmt" "io" "os" + "strings" "infini.sh/agent/lib/reader" "infini.sh/agent/lib/reader/linenumber" @@ -23,34 +25,60 @@ type Harvester struct { file *os.File config Config offset int64 + // gzipWrapped: the file is a rotated .gz archive — immutable, always + // read from the start of the DECOMPRESSED stream (offsets into it are + // not comparable across runs, so state tracking is mtime-based only). + gzipWrapped bool + // src is the reading source (raw file or gzip-decompressed stream). + src io.ReadCloser encodingFactory encoding.EncodingFactory encoding encoding.Encoding } +// IsGzipFile reports whether the path looks like a gzip archive — the +// detector routes rotated *.gz logs here so compressed history is ingested. +func IsGzipFile(path string) bool { + return strings.HasSuffix(strings.ToLower(path), ".gz") +} + func NewHarvester(path string, offset int64) (*Harvester, error) { f, err := readOpen(path) if f == nil || err != nil { return nil, errors.Errorf("failed to open file(%s),%v", path, err) } - _, err = f.Seek(offset, io.SeekStart) - if err != nil { - return nil, err - } h := &Harvester{ file: f, config: defaultConfig(), offset: offset, } + // Reading source: the raw file, or its decompressed stream for rotated + // .gz archives (immutable — always from the start, offsets in the + // decompressed stream are not comparable across runs). + src := io.ReadCloser(f) + if IsGzipFile(path) { + zr, zerr := gzip.NewReader(f) + if zerr != nil { + _ = f.Close() + return nil, errors.Errorf("failed to open gzip file(%s),%v", path, zerr) + } + src = zr + h.gzipWrapped = true + h.offset = 0 + } else if _, err := f.Seek(offset, io.SeekStart); err != nil { + _ = f.Close() + return nil, err + } encodingFactory, ok := encoding.FindEncoding(h.config.Encoding) if !ok || encodingFactory == nil { return nil, fmt.Errorf("unknown encoding('%v')", h.config.Encoding) } h.encodingFactory = encodingFactory - h.encoding, err = h.encodingFactory(f) + h.encoding, err = h.encodingFactory(src) if err != nil { return nil, err } + h.src = src return h, nil } @@ -68,7 +96,7 @@ func (h *Harvester) NewJsonFileReader(pattern string, showLineNumber bool) (read } encReaderMaxBytes := h.config.MaxBytes * 4 - r, err = readfile.NewEncodeReader(h.file, readfile.Config{ + r, err = readfile.NewEncodeReader(h.src, readfile.Config{ Codec: h.encoding, BufferSize: h.config.BufferSize, MaxBytes: encReaderMaxBytes, @@ -107,7 +135,7 @@ func (h *Harvester) NewLogFileReader(pattern string, showLineNumber bool) (reade return nil, fmt.Errorf("file is nil") } encReaderMaxBytes := h.config.MaxBytes * 4 - r, err = readfile.NewEncodeReader(h.file, readfile.Config{ + r, err = readfile.NewEncodeReader(h.src, readfile.Config{ Codec: h.encoding, BufferSize: h.config.BufferSize, MaxBytes: encReaderMaxBytes, diff --git a/lib/util/saferead.go b/lib/util/saferead.go new file mode 100644 index 0000000..654124c --- /dev/null +++ b/lib/util/saferead.go @@ -0,0 +1,268 @@ +/* Copyright © INFINI Ltd. All rights reserved. + * Web: https://infinilabs.com + * Email: hello#infini.ltd */ + +// Vendored copy of the framework's core/util saferead (framework PR #420, +// instance-log-tail-api). Keep in sync with upstream and switch callers +// back to infini.sh/framework/core/util once the framework release that +// ships it is a hard dependency here. + +package util + +import ( + "fmt" + "os" + "path/filepath" + "strings" +) + +// Read-side file access control: any handler that serves file contents +// driven by external input (a caller-supplied directory, a file name from +// a query string) must resolve the path through a ReadGuard so the +// endpoint cannot be turned into a general-purpose file reader. +// +// Two layers, deny always wins: +// 1. IsSystemReadPath — OS-critical trees are never readable, even when +// a whitelist root was misconfigured to include them. +// 2. ReadGuard — files must live strictly below one of the explicitly +// allowed roots (symlink-resolved on both sides). + +// systemReadPrefixes lists trees that never legitimately hold readable +// service logs, so serving reads from them is refused outright. It is +// deliberately narrower than safedelete's neverTouch: packaged products +// keep logs under /usr/share//logs (RPM/DEB installs) and +// /var/log, /opt, /home are standard log locations — those stay +// reachable, constrained by the whitelist instead. +var systemReadPrefixes = []string{ + "/bin", "/sbin", "/boot", "/dev", "/etc", "/proc", "/sys", "/run", + "/lib", "/lib64", "/libx32", "/root", "/kernel", + "/System", "/private/etc", + "/usr/bin", "/usr/sbin", "/usr/lib", "/usr/lib64", "/usr/libx32", +} + +// windowsSystemReadPrefixes are matched against the volume-relative part +// of a drive path (C:/Windows/... -> /Windows/...) case-insensitively. +var windowsSystemReadPrefixes = []string{ + "/Windows", "/Program Files/WindowsApps", +} + +// tooBroadReadRoots lists top-level trees that hold far more than logs; +// whitelisting one of them verbatim (instead of a specific subdirectory) +// is treated as a misconfiguration and rejected. Their contents remain +// reachable through narrower roots. +var tooBroadReadRoots = []string{ + "/usr", "/var", "/opt", "/srv", "/home", "/mnt", "/media", "/tmp", + "/private/var", "/private/tmp", +} + +// IsSystemReadPath reports whether p is an OS-critical location that must +// never be served by a file-reading API: the filesystem root, a +// systemReadPrefixes member or anything under it (drive roots and the +// listed Windows trees on Windows). +func IsSystemReadPath(p string) bool { + return isSystemReadCanonical(canonicalReadPath(p)) +} + +// canonicalPath absolutizes and symlink-resolves p. When p itself does not +// exist yet, the deepest EXISTING ancestor is resolved instead, so a not- +// yet-created target under a symlinked root canonicalizes identically to +// the root — and a link jumping out of the allowed root cannot smuggle a +// target past the containment check either way. (Vendored from the +// framework's core/util safedelete.go.) +func canonicalPath(p string) string { + abs, err := filepath.Abs(p) + if err != nil { + abs = filepath.Clean(p) + } + if resolved, err := filepath.EvalSymlinks(abs); err == nil { + return resolved + } + rest := "" + dir := abs + for { + parent := filepath.Dir(dir) + if parent == dir { // reached "/" + break + } + rest = filepath.Join(filepath.Base(dir), rest) + if resolved, err := filepath.EvalSymlinks(parent); err == nil { + return filepath.Join(resolved, rest) + } + dir = parent + } + return abs +} + +// canonicalReadPath canonicalizes p and maps macOS firmlink paths +// (/System/Volumes/Data/...) back to their firmlink view (/home, /opt, +// ...) so every guard layer compares the same spelling. +func canonicalReadPath(p string) string { + c := canonicalPath(p) + if strings.HasPrefix(c, "/System/Volumes/Data/") { + c = c[len("/System/Volumes/Data"):] + } + return c +} + +// isSystemReadCanonical applies the deny checks to an already-canonical +// path so the Windows rules stay testable on every platform. +func isSystemReadCanonical(c string) bool { + if c == "/" || c == string(filepath.Separator) { + return true + } + // drive-qualified path (Windows): compare the volume-relative part. + // The separator check keeps unix paths whose second character happens + // to be ':' (e.g. /t:mp) on the unix branch. + if len(c) > 2 && c[1] == ':' && (c[2] == '/' || c[2] == '\\') { + c = filepath.ToSlash(c) + if len(c) <= 3 { // bare drive root, e.g. C:/ + return true + } + rel := strings.ToLower(c[2:]) // /windows/system32 + for _, prefix := range windowsSystemReadPrefixes { + prefix = strings.ToLower(prefix) + if rel == prefix || strings.HasPrefix(rel, prefix+"/") { + return true + } + } + return false + } + for _, prefix := range systemReadPrefixes { + if c == prefix || strings.HasPrefix(c, prefix+string(filepath.Separator)) { + return true + } + } + return false +} + +// isTooBroadReadRoot reports whether an already-canonical directory is a +// wholesale top-level tree rather than a specific log location. +func isTooBroadReadRoot(c string) bool { + for _, broad := range tooBroadReadRoots { + if c == broad { + return true + } + } + return false +} + +// ReadGuard confines file reads to a fixed set of allowed roots. Roots +// are canonicalized (absolute + symlink-resolved) once at construction; +// every later ResolveUnder call re-resolves the target the same way so a +// symlink cannot smuggle a path outside the roots. +type ReadGuard struct { + roots []string +} + +// NewReadGuard validates and canonicalizes the allowed roots. Every root +// must exist, be a directory, be a specific-enough location (not a +// wholesale /usr or /var), and not be a system path — a bad root fails +// loudly instead of being silently dropped. +func NewReadGuard(roots ...string) (*ReadGuard, error) { + g := &ReadGuard{} + for _, root := range roots { + root = strings.TrimSpace(root) + if root == "" { + continue + } + c := canonicalReadPath(root) + if IsSystemReadPath(c) { + return nil, fmt.Errorf("refusing to serve reads from system path: %s", c) + } + if isTooBroadReadRoot(c) { + return nil, fmt.Errorf("allowed path %s is too broad, whitelist a specific log directory instead", c) + } + fi, err := os.Stat(c) + if err != nil { + return nil, fmt.Errorf("allowed path %s is not accessible: %v", c, err) + } + if !fi.IsDir() { + return nil, fmt.Errorf("allowed path %s is not a directory", c) + } + g.roots = append(g.roots, c) + } + if len(g.roots) == 0 { + return nil, fmt.Errorf("no allowed read paths configured") + } + return g, nil +} + +// Roots returns the canonical allowed roots. +func (g *ReadGuard) Roots() []string { + return append([]string(nil), g.roots...) +} + +// Contains reports whether path canonicalizes to exactly one of the +// allowed roots. +func (g *ReadGuard) Contains(path string) bool { + if strings.TrimSpace(path) == "" { + return false + } + c := canonicalReadPath(path) + for _, r := range g.roots { + if c == r { + return true + } + } + return false +} + +// ContainsUnder reports whether path canonicalizes to one of the allowed +// roots or to a location strictly below one of them — the check for +// caller-supplied base directories that may be a subdirectory of a +// whitelisted root (e.g. a GC log dir configured under path.logs). +func (g *ReadGuard) ContainsUnder(path string) bool { + if strings.TrimSpace(path) == "" { + return false + } + c := canonicalReadPath(path) + for _, r := range g.roots { + if c == r || strings.HasPrefix(c, r+string(filepath.Separator)) { + return true + } + } + return false +} + +// ResolveUnder resolves file against base (which must be an allowed root +// or live below one) and returns its canonical path when it lands strictly +// below an allowed root, outside every system tree, and is an existing +// regular file (so devices, fifos and sockets can never be served). +// file may be relative (joined to base) or absolute (accepted only when +// it still resolves inside the roots); the roots — not base — remain the +// trust boundary, so a file name escaping base but staying inside the +// root is still served. +func (g *ReadGuard) ResolveUnder(base, file string) (string, error) { + if strings.TrimSpace(file) == "" { + return "", fmt.Errorf("file is required") + } + if !g.ContainsUnder(base) { + return "", fmt.Errorf("path [%v] is not an allowed directory", base) + } + target := file + if !filepath.IsAbs(target) { + target = filepath.Join(canonicalReadPath(base), target) + } + c := canonicalReadPath(target) + if IsSystemReadPath(c) { + return "", fmt.Errorf("path [%v] is a system path", file) + } + contained := false + for _, r := range g.roots { + if strings.HasPrefix(c, r+string(filepath.Separator)) { + contained = true + break + } + } + if !contained { + return "", fmt.Errorf("file [%v] is outside of the allowed directories", file) + } + fi, err := os.Stat(c) + if err != nil { + return "", fmt.Errorf("cannot access file [%v]: %v", file, err) + } + if !fi.Mode().IsRegular() { + return "", fmt.Errorf("[%v] is not a regular file", file) + } + return c, nil +} diff --git a/lib/util/saferead_test.go b/lib/util/saferead_test.go new file mode 100644 index 0000000..35a75bc --- /dev/null +++ b/lib/util/saferead_test.go @@ -0,0 +1,276 @@ +/* Copyright © INFINI Ltd. All rights reserved. + * Web: https://infinilabs.com + * Email: hello#infini.ltd */ + +package util + +import ( + "os" + "path/filepath" + "strings" + "testing" +) + +func TestIsSystemReadPath(t *testing.T) { + denied := []string{ + "/", + "/etc", "/etc/passwd", "/etc/ssl/private/server.key", + "/bin/ls", "/sbin", "/boot/grub", + "/usr/bin/ls", "/usr/sbin/x", "/usr/lib/x.so", "/usr/lib64/x.so", + "/root/.ssh/id_rsa", "/kernel", + "/System/Library", "/private/etc/hosts", + "/proc/self/environ", "/sys/kernel", "/dev/sda", "/run/secrets/token", + // traversal into a system tree + "/var/log/../../../etc/passwd", + } + for _, p := range denied { + if !IsSystemReadPath(p) { + t.Errorf("expected [%s] to be a system read path", p) + } + } + + allowed := []string{ + "/var/log/elasticsearch", + "/usr/share/easysearch/logs/server.log", + "/opt/myapp/logs", + "/home/user/logs/app.log", + "/srv/service/logs", + } + for _, p := range allowed { + if IsSystemReadPath(p) { + t.Errorf("expected [%s] to be a readable location", p) + } + } +} + +func TestIsSystemReadPathUnixPathWithColon(t *testing.T) { + // second-char-colon unix paths must stay on the unix branch: not + // system paths themselves, and real system paths still denied + if IsSystemReadPath("/t:mp/logs/server.log") { + t.Error("expected colon-containing unix path to be readable") + } + if !IsSystemReadPath("/etc/passwd") { + t.Error("expected /etc/passwd to stay denied") + } +} + +func TestIsSystemReadCanonicalWindows(t *testing.T) { + denied := []string{ + "C:/Windows", "C:/Windows/System32/cmd.exe", "c:/windows/system32/x.dll", + "C:/Program Files/WindowsApps/App", + } + for _, p := range denied { + if !isSystemReadCanonical(p) { + t.Errorf("expected [%s] to be a system read path", p) + } + } + allowed := []string{ + "C:/Program Files/Elasticsearch/logs", + "C:/elastic/logs/server.log", + "C:/ProgramData/myapp/logs", + } + for _, p := range allowed { + if isSystemReadCanonical(p) { + t.Errorf("expected [%s] to be a readable location", p) + } + } +} + +func TestNewReadGuardRootValidation(t *testing.T) { + dir := t.TempDir() + + if _, err := NewReadGuard(dir); err != nil { + t.Fatalf("expected temp dir to be a valid root: %v", err) + } + + for _, bad := range []string{"/etc", "/usr/bin", "/root", "/"} { + if _, err := NewReadGuard(bad); err == nil { + t.Errorf("expected system path [%s] to be rejected as root", bad) + } + } + for _, broad := range []string{"/usr", "/var", "/opt", "/home", "/tmp"} { + if _, err := NewReadGuard(broad); err == nil { + t.Errorf("expected broad root [%s] to be rejected", broad) + } + } + + file := filepath.Join(dir, "notadir.log") + if err := os.WriteFile(file, []byte("x"), 0644); err != nil { + t.Fatal(err) + } + if _, err := NewReadGuard(file); err == nil { + t.Error("expected non-directory root to be rejected") + } + missing := filepath.Join(dir, "does-not-exist") + if _, err := NewReadGuard(missing); err == nil { + t.Error("expected missing root to be rejected") + } + if _, err := NewReadGuard(); err == nil { + t.Error("expected empty roots to be rejected") + } +} + +func TestReadGuardResolveUnder(t *testing.T) { + dir := t.TempDir() + guard, err := NewReadGuard(dir) + if err != nil { + t.Fatal(err) + } + + logFile := canonicalPath(filepath.Join(dir, "server.log")) + if err := os.WriteFile(logFile, []byte("line"), 0644); err != nil { + t.Fatal(err) + } + + if got, err := guard.ResolveUnder(dir, "server.log"); err != nil || got != logFile { + t.Errorf("relative resolve failed: %v %v", got, err) + } + if got, err := guard.ResolveUnder(dir, logFile); err != nil || got != logFile { + t.Errorf("absolute-inside resolve failed: %v %v", got, err) + } + + for _, bad := range []string{ + "../../../etc/passwd", + "/etc/passwd", + "", + ".", + "missing.log", + } { + if _, err := guard.ResolveUnder(dir, bad); err == nil { + t.Errorf("expected file [%q] to be rejected", bad) + } + } + + if _, err := guard.ResolveUnder("/etc", "passwd"); err == nil { + t.Error("expected base outside the roots to be rejected") + } +} + +func TestReadGuardRejectsNonRegularFiles(t *testing.T) { + dir := t.TempDir() + guard, err := NewReadGuard(dir) + if err != nil { + t.Fatal(err) + } + + sub := filepath.Join(dir, "nested") + if err := os.Mkdir(sub, 0755); err != nil { + t.Fatal(err) + } + if _, err := guard.ResolveUnder(dir, "nested"); err == nil { + t.Error("expected directory passed as file to be rejected") + } + + fifo := filepath.Join(dir, "pipe.log") + if err := mkFifo(fifo); err != nil { + t.Skipf("cannot create fifo on this platform: %v", err) + } + if _, err := guard.ResolveUnder(dir, "pipe.log"); err == nil { + t.Error("expected fifo to be rejected") + } +} + +func TestReadGuardContainsUnder(t *testing.T) { + dir := t.TempDir() + sub := filepath.Join(dir, "gc") + if err := os.Mkdir(sub, 0755); err != nil { + t.Fatal(err) + } + guard, err := NewReadGuard(dir) + if err != nil { + t.Fatal(err) + } + + if !guard.ContainsUnder(dir) { + t.Error("expected root itself to match") + } + if !guard.ContainsUnder(sub) { + t.Error("expected subdirectory of the root to match") + } + if guard.ContainsUnder(dir + "-sibling") { + t.Error("expected prefix sibling to be rejected") + } + if guard.ContainsUnder("/etc") { + t.Error("expected outside path to be rejected") + } + + // base may be a subdirectory; the root stays the trust boundary + logFile := filepath.Join(dir, "server.log") + if err := os.WriteFile(logFile, []byte("line"), 0644); err != nil { + t.Fatal(err) + } + if got, err := guard.ResolveUnder(sub, "../server.log"); err != nil || got != canonicalPath(logFile) { + t.Errorf("expected escape from base but not root to resolve, got %v %v", got, err) + } + if _, err := guard.ResolveUnder(sub, "../../../etc/passwd"); err == nil { + t.Error("expected escape out of the root to be rejected") + } +} + +func TestReadGuardSymlinkEscape(t *testing.T) { + dir := t.TempDir() + guard, err := NewReadGuard(dir) + if err != nil { + t.Fatal(err) + } + + logFile := canonicalPath(filepath.Join(dir, "server.log")) + if err := os.WriteFile(logFile, []byte("line"), 0644); err != nil { + t.Fatal(err) + } + + link := filepath.Join(dir, "etc-link") + if err := os.Symlink("/etc", link); err != nil { + t.Fatal(err) + } + if _, err := guard.ResolveUnder(dir, "etc-link/passwd"); err == nil { + t.Error("expected symlink escape into /etc to be rejected") + } + + // a root reached through its own symlink still canonicalizes to itself + linkToRoot := dir + "-via-link" + if err := os.Symlink(dir, linkToRoot); err != nil { + t.Fatal(err) + } + if !guard.Contains(linkToRoot) { + t.Error("expected symlinked root variant to be recognized") + } + if got, err := guard.ResolveUnder(linkToRoot, "server.log"); err != nil || got != logFile { + t.Errorf("resolve via symlinked root failed: %v %v", got, err) + } +} + +func TestReadGuardPrefixSafety(t *testing.T) { + base := t.TempDir() + sibling, err := os.MkdirTemp(filepath.Dir(base), filepath.Base(base)+"-sibling") + if err != nil { + t.Fatal(err) + } + defer os.RemoveAll(sibling) + if strings.HasPrefix(sibling, base+string(filepath.Separator)) { + t.Fatalf("test requires sibling outside base: %s vs %s", sibling, base) + } + + guard, err := NewReadGuard(base) + if err != nil { + t.Fatal(err) + } + escape := filepath.Join(sibling, "secret.log") + if err := os.WriteFile(escape, []byte("secret"), 0644); err != nil { + t.Fatal(err) + } + if _, err := guard.ResolveUnder(base, escape); err == nil { + t.Error("expected prefix-sibling path (base vs base-sibling) to be rejected") + } +} + +func TestIsTooBroadReadRoot(t *testing.T) { + for _, broad := range []string{"/usr", "/var", "/private/var"} { + if !isTooBroadReadRoot(broad) { + t.Errorf("expected [%s] to be too broad", broad) + } + } + if isTooBroadReadRoot("/var/log/myapp") { + t.Error("expected specific subdir to be acceptable") + } +} diff --git a/lib/util/saferead_test_unix.go b/lib/util/saferead_test_unix.go new file mode 100644 index 0000000..378477c --- /dev/null +++ b/lib/util/saferead_test_unix.go @@ -0,0 +1,13 @@ +//go:build !windows + +/* Copyright © INFINI Ltd. All rights reserved. + * Web: https://infinilabs.com + * Email: hello#infini.ltd */ + +package util + +import "syscall" + +func mkFifo(path string) error { + return syscall.Mkfifo(path, 0600) +} diff --git a/lib/util/saferead_test_windows.go b/lib/util/saferead_test_windows.go new file mode 100644 index 0000000..a2e0263 --- /dev/null +++ b/lib/util/saferead_test_windows.go @@ -0,0 +1,13 @@ +//go:build windows + +/* Copyright © INFINI Ltd. All rights reserved. + * Web: https://infinilabs.com + * Email: hello#infini.ltd */ + +package util + +import "errors" + +func mkFifo(path string) error { + return errors.New("fifo is not supported on windows") +} diff --git a/main.go b/main.go index 9f9e6eb..0818fef 100755 --- a/main.go +++ b/main.go @@ -8,6 +8,7 @@ import ( _ "expvar" public "infini.sh/agent/.public" "infini.sh/agent/config" + _ "infini.sh/agent/internal/managed" _ "infini.sh/agent/plugin" api3 "infini.sh/agent/plugin/api" "infini.sh/framework" @@ -19,10 +20,12 @@ import ( "infini.sh/framework/core/util" "infini.sh/framework/core/vfs" "infini.sh/framework/modules/api" + _ "infini.sh/framework/modules/configs/reverseclient" "infini.sh/framework/modules/elastic" "infini.sh/framework/modules/keystore" "infini.sh/framework/modules/metrics" "infini.sh/framework/modules/pipeline" + _ "infini.sh/framework/modules/queue" queue2 "infini.sh/framework/modules/queue/disk_queue" "infini.sh/framework/modules/security" stats2 "infini.sh/framework/modules/stats" @@ -32,6 +35,9 @@ import ( _ "infini.sh/framework/plugins/elastic/indexing_merge" _ "infini.sh/framework/plugins/http" _ "infini.sh/framework/plugins/queue/consumer" + // Kafka bus queue backend (agent-side produce only: logs_processor + // writes straight into a kafka-typed queue) + _ "infini.sh/framework/plugins/queue/kafka_queue" "infini.sh/framework/plugins/simple_kv" "os" "runtime" diff --git a/plugin/api/log.go b/plugin/api/log.go index 1ed404b..4572d32 100644 --- a/plugin/api/log.go +++ b/plugin/api/log.go @@ -30,11 +30,27 @@ func (handler *AgentAPI) getSearchLogFiles(w http.ResponseWriter, req *http.Requ handler.WriteJSON(w, err.Error(), http.StatusInternalServerError) return } + guard, err := esLogsReadGuard() + if err != nil { + handler.WriteError(w, err.Error(), http.StatusForbidden) + return + } logsPaths := normalizeJSONLogsPaths(reqBody.LogsPath) if len(logsPaths) == 0 { handler.WriteError(w, "miss param logs_path", http.StatusInternalServerError) return } + var denied []string + for _, logsPath := range logsPaths { + if !guard.ContainsUnder(logsPath) { + denied = append(denied, logsPath) + } + } + if len(denied) > 0 { + log.Warnf("rejected search log files request outside the whitelist: %v", denied) + handler.WriteError(w, fmt.Sprintf("logs_path %v is not allowed, configure elasticsearch_logs.allowed_paths or check elasticsearch discovery", denied), http.StatusForbidden) + return + } var files []util.MapStr var errors []string @@ -53,6 +69,9 @@ func (handler *AgentAPI) getSearchLogFiles(w http.ResponseWriter, req *http.Requ if info.IsDir() { continue } + if !info.Type().IsRegular() { + continue + } fInfo, err := info.Info() if err != nil { appendError("failed to read file info in logs directory [%s], file=[%s]: %v", logsPath, info.Name(), err) @@ -97,10 +116,22 @@ func (handler *AgentAPI) readSearchLogFile(w http.ResponseWriter, req *http.Requ return } - logFilePath, err := safeJoinLogsFile(reqBody.LogsPath, reqBody.FileName) + guard, err := esLogsReadGuard() + if err != nil { + handler.WriteError(w, err.Error(), http.StatusForbidden) + return + } + if !guard.ContainsUnder(reqBody.LogsPath) { + log.Warnf("rejected search log read request outside the whitelist, logs_path=[%s]", reqBody.LogsPath) + handler.WriteError(w, fmt.Sprintf("logs_path [%s] is not allowed, configure elasticsearch_logs.allowed_paths or check elasticsearch discovery", reqBody.LogsPath), http.StatusForbidden) + return + } + // resolves inside the whitelisted dir only: traversal, symlink escapes + // and non-regular files (devices, fifos) are rejected here + logFilePath, err := guard.ResolveUnder(reqBody.LogsPath, reqBody.FileName) if err != nil { log.Errorf("invalid search log file request, logs_path=[%s], file_name=[%s]: %v", reqBody.LogsPath, reqBody.FileName, err) - handler.WriteJSON(w, err.Error(), http.StatusInternalServerError) + handler.WriteError(w, err.Error(), http.StatusBadRequest) return } if reqBody.StartLineNumber < 0 { @@ -115,14 +146,14 @@ func (handler *AgentAPI) readSearchLogFile(w http.ResponseWriter, req *http.Requ err = os.MkdirAll(fileDir, os.ModePerm) if err != nil { log.Errorf("failed to create temporary log directory [%s] for source [%s]: %v", fileDir, logFilePath, err) - handler.WriteJSON(w, err.Error(), http.StatusInternalServerError) + handler.WriteError(w, err.Error(), http.StatusInternalServerError) return } } err = agent_util.UnpackGzipFile(logFilePath, tmpFilePath) if err != nil { log.Errorf("failed to unpack gzip log file from [%s] to [%s]: %v", logFilePath, tmpFilePath, err) - handler.WriteJSON(w, err.Error(), http.StatusInternalServerError) + handler.WriteError(w, err.Error(), http.StatusInternalServerError) return } } @@ -224,28 +255,3 @@ func normalizeJSONLogsPaths(raw interface{}) []string { } return result } - -func safeJoinLogsFile(logsPath, fileName string) (string, error) { - logsPath = strings.TrimSpace(logsPath) - fileName = strings.TrimSpace(fileName) - if logsPath == "" || fileName == "" { - return "", fmt.Errorf("invalid log file request") - } - - expanded, err := agent_util.ExpandHomeDir(logsPath) - if err != nil { - return "", err - } - logsPath = expanded - - basePath := filepath.Clean(logsPath) - fullPath := filepath.Clean(filepath.Join(basePath, fileName)) - rel, err := filepath.Rel(basePath, fullPath) - if err != nil { - return "", err - } - if rel == ".." || strings.HasPrefix(rel, ".."+string(filepath.Separator)) { - return "", fmt.Errorf("invalid log file path") - } - return fullPath, nil -} diff --git a/plugin/api/log_test.go b/plugin/api/log_test.go index 4bfe113..7ea12a8 100644 --- a/plugin/api/log_test.go +++ b/plugin/api/log_test.go @@ -1,6 +1,17 @@ package api -import "testing" +import ( + "bytes" + "net/http/httptest" + "os" + "path/filepath" + "strings" + "testing" + + httprouter "infini.sh/framework/core/api/router" + "infini.sh/framework/core/elastic" + "infini.sh/framework/core/util" +) func TestNormalizeJSONLogsPathsSupportsStringAndArray(t *testing.T) { paths := normalizeJSONLogsPaths("/var/log/elasticsearch") @@ -14,13 +25,260 @@ func TestNormalizeJSONLogsPathsSupportsStringAndArray(t *testing.T) { } } -func TestSafeJoinLogsFileRejectsTraversal(t *testing.T) { - path, err := safeJoinLogsFile("/var/log/elasticsearch", "server.log") - if err != nil || path != "/var/log/elasticsearch/server.log" { - t.Fatalf("expected normal path, got path=%q err=%v", path, err) +func withESLogWhitelist(t *testing.T, dirs ...string) { + t.Helper() + orig := esLogWhitelistLoader + esLogWhitelistLoader = func() ([]string, error) { return dirs, nil } + resetESLogWhitelistCache() + t.Cleanup(func() { + esLogWhitelistLoader = orig + resetESLogWhitelistCache() + }) +} + +func TestGetSearchLogFilesRejectsPathOutsideWhitelist(t *testing.T) { + dir := t.TempDir() + withESLogWhitelist(t, dir) + + handler := AgentAPI{} + req := httptest.NewRequest("POST", "/elasticsearch/logs/_list", strings.NewReader(`{"logs_path":"/etc"}`)) + w := httptest.NewRecorder() + handler.getSearchLogFiles(w, req, httprouter.Params{}) + if w.Code != 403 { + t.Fatalf("expected 403, got %d: %s", w.Code, w.Body.String()) + } +} + +func TestGetSearchLogFilesEmptyWhitelistDenied(t *testing.T) { + withESLogWhitelist(t) + + handler := AgentAPI{} + req := httptest.NewRequest("POST", "/elasticsearch/logs/_list", strings.NewReader(`{"logs_path":"/var/log/elasticsearch"}`)) + w := httptest.NewRecorder() + handler.getSearchLogFiles(w, req, httprouter.Params{}) + if w.Code != 403 { + t.Fatalf("expected 403 when no whitelist can be established, got %d: %s", w.Code, w.Body.String()) + } +} + +func TestGetSearchLogFilesListsAllowedPath(t *testing.T) { + dir := t.TempDir() + withESLogWhitelist(t, dir) + + if err := os.WriteFile(filepath.Join(dir, "server.log"), []byte("line1\n"), 0644); err != nil { + t.Fatal(err) + } + if err := os.WriteFile(filepath.Join(dir, "notes.txt"), []byte("x"), 0644); err != nil { + t.Fatal(err) + } + + handler := AgentAPI{} + body := `{"logs_path":` + quoteJSON(dir) + `}` + req := httptest.NewRequest("POST", "/elasticsearch/logs/_list", strings.NewReader(body)) + w := httptest.NewRecorder() + handler.getSearchLogFiles(w, req, httprouter.Params{}) + if w.Code != 200 { + t.Fatalf("expected 200, got %d: %s", w.Code, w.Body.String()) + } + resp := w.Body.String() + if !strings.Contains(resp, `"success":true`) || !strings.Contains(resp, "server.log") { + t.Errorf("expected server.log in successful listing: %s", resp) + } + if !strings.Contains(resp, "notes.txt") { + t.Errorf("expected listing to keep every regular file: %s", resp) + } +} + +func TestReadSearchLogFileWhitelistAndTraversal(t *testing.T) { + dir := t.TempDir() + withESLogWhitelist(t, dir) + + logFile := filepath.Join(dir, "server.log") + if err := os.WriteFile(logFile, []byte("line1\nline2\n"), 0644); err != nil { + t.Fatal(err) + } + + handler := AgentAPI{} + + // base outside the whitelist + req := httptest.NewRequest("POST", "/elasticsearch/logs/_read", + strings.NewReader(`{"logs_path":"/etc","file_name":"passwd","lines":10}`)) + w := httptest.NewRecorder() + handler.readSearchLogFile(w, req, httprouter.Params{}) + if w.Code != 403 { + t.Fatalf("expected 403 for logs_path outside whitelist, got %d: %s", w.Code, w.Body.String()) + } + + // traversal within an allowed base + req = httptest.NewRequest("POST", "/elasticsearch/logs/_read", + strings.NewReader(`{"logs_path":`+quoteJSON(dir)+`,"file_name":"../secret.log","lines":10}`)) + w = httptest.NewRecorder() + handler.readSearchLogFile(w, req, httprouter.Params{}) + if w.Code != 400 { + t.Fatalf("expected 400 for traversal, got %d: %s", w.Code, w.Body.String()) + } + + // symlink escape out of the allowed base + link := filepath.Join(dir, "escape.log") + if err := os.Symlink("/etc/passwd", link); err != nil { + t.Fatal(err) + } + req = httptest.NewRequest("POST", "/elasticsearch/logs/_read", + strings.NewReader(`{"logs_path":`+quoteJSON(dir)+`,"file_name":"escape.log","lines":10}`)) + w = httptest.NewRecorder() + handler.readSearchLogFile(w, req, httprouter.Params{}) + if w.Code != 400 { + t.Fatalf("expected 400 for symlink escape, got %d: %s", w.Code, w.Body.String()) + } + + // normal read inside the whitelist + req = httptest.NewRequest("POST", "/elasticsearch/logs/_read", + strings.NewReader(`{"logs_path":`+quoteJSON(dir)+`,"file_name":"server.log","lines":10,"start_line_number":1}`)) + w = httptest.NewRecorder() + handler.readSearchLogFile(w, req, httprouter.Params{}) + if w.Code != 200 { + t.Fatalf("expected 200, got %d: %s", w.Code, w.Body.String()) + } + if !strings.Contains(w.Body.String(), "line2") { + t.Errorf("expected file content in response: %s", w.Body.String()) + } +} + +func TestNodeLogDirs(t *testing.T) { + // path.logs as string (home/logs also kept as a candidate) + info := &elastic.NodesInfo{Settings: map[string]interface{}{ + "path": map[string]interface{}{ + "logs": "/var/log/easysearch", + "home": "/usr/share/easysearch", + }, + }} + if dirs := nodeLogDirs(info); len(dirs) != 2 || dirs[0] != "/var/log/easysearch" || dirs[1] != filepath.Join("/usr/share/easysearch", "logs") { + t.Fatalf("expected logs and home/logs dirs, got %#v", dirs) + } + + // path.logs as array (multi-path) + info.Settings = map[string]interface{}{ + "path": map[string]interface{}{"logs": []interface{}{"/data1/logs", "/data2/logs"}}, + } + if dirs := nodeLogDirs(info); len(dirs) != 2 { + t.Fatalf("expected multi-path logs dirs, got %#v", dirs) } - if _, err := safeJoinLogsFile("/var/log/elasticsearch", "../server.log"); err == nil { - t.Fatal("expected traversal to be rejected") + // fallback to path.home/logs + info.Settings = map[string]interface{}{ + "path": map[string]interface{}{"home": "/usr/share/easysearch"}, } + if dirs := nodeLogDirs(info); len(dirs) != 1 || dirs[0] != filepath.Join("/usr/share/easysearch", "logs") { + t.Fatalf("expected home/logs fallback, got %#v", dirs) + } + + if dirs := nodeLogDirs(nil); dirs != nil { + t.Fatalf("expected no dirs for nil node info, got %#v", dirs) + } +} + +func TestGetSearchLogFilesAllowsSubdirOfWhitelistedRoot(t *testing.T) { + dir := t.TempDir() + sub := filepath.Join(dir, "gc") + if err := os.Mkdir(sub, 0755); err != nil { + t.Fatal(err) + } + withESLogWhitelist(t, dir) + + if err := os.WriteFile(filepath.Join(sub, "gc.log"), []byte("gc\n"), 0644); err != nil { + t.Fatal(err) + } + + handler := AgentAPI{} + req := httptest.NewRequest("POST", "/elasticsearch/logs/_list", strings.NewReader(`{"logs_path":`+quoteJSON(sub)+`}`)) + w := httptest.NewRecorder() + handler.getSearchLogFiles(w, req, httprouter.Params{}) + if w.Code != 200 { + t.Fatalf("expected 200 for subdirectory of whitelisted root, got %d: %s", w.Code, w.Body.String()) + } + if !strings.Contains(w.Body.String(), "gc.log") { + t.Errorf("expected gc.log in listing: %s", w.Body.String()) + } +} + +func TestCmdlineLogDirs(t *testing.T) { + // -Des.path.logs override + dirs := cmdlineLogDirs(`java -Des.path.home=/usr/share/easysearch -Des.path.logs=/var/log/easysearch -Xlog:gc*:file=/var/log/easysearch/gc.log:uptime,tags`) + if len(dirs) != 2 || dirs[0] != "/var/log/easysearch" || dirs[1] != "/var/log/easysearch" { + t.Fatalf("expected cmdline logs dir and gc dir, got %#v", dirs) + } + + // gc file outside path.logs must be whitelisted too + dirs = cmdlineLogDirs(`java -Des.path.home=/usr/share/easysearch -Des.path.logs=/var/log/easysearch -Xlog:gc*:file=/var/log/gc/es_gc.log:uptime`) + if len(dirs) != 2 || dirs[1] != "/var/log/gc" { + t.Fatalf("expected gc dir outside path.logs, got %#v", dirs) + } + + // no explicit path.logs: derive home/logs + dirs = cmdlineLogDirs(`java -Des.path.home=/opt/easysearch`) + if len(dirs) != 1 || dirs[0] != filepath.Join("/opt/easysearch", "logs") { + t.Fatalf("expected home/logs dir, got %#v", dirs) + } + + // quoted relative gc file resolved against home + dirs = cmdlineLogDirs(`java -Des.path.home=/opt/es -Xlog:file="logs/gc.log"`) + if len(dirs) != 2 || dirs[1] != filepath.Join("/opt/es", "logs") { + t.Fatalf("expected relative gc file joined on home, got %#v", dirs) + } +} + +func TestBuildESLogsReadGuardSkipsInvalidRoots(t *testing.T) { + dir := t.TempDir() + orig := esLogWhitelistLoader + esLogWhitelistLoader = func() ([]string, error) { + return []string{"/etc", filepath.Join(dir, "missing"), dir}, nil + } + resetESLogWhitelistCache() + defer func() { + esLogWhitelistLoader = orig + resetESLogWhitelistCache() + }() + + guard, err := esLogsReadGuard() + if err != nil { + t.Fatalf("expected valid roots to survive: %v", err) + } + if !guard.Contains(dir) { + t.Error("expected the valid root to be whitelisted") + } + if guard.Contains("/etc") { + t.Error("expected system path to be excluded from the whitelist") + } +} + +func TestReadSearchLogFileResponseShape(t *testing.T) { + dir := t.TempDir() + withESLogWhitelist(t, dir) + + if err := os.WriteFile(filepath.Join(dir, "server.log"), []byte("hello\n"), 0644); err != nil { + t.Fatal(err) + } + + handler := AgentAPI{} + var buf bytes.Buffer + buf.WriteString(`{"logs_path":"`) + buf.WriteString(strings.ReplaceAll(dir, `\`, `\\`)) + buf.WriteString(`","file_name":"server.log","lines":5,"start_line_number":0}`) + req := httptest.NewRequest("POST", "/elasticsearch/logs/_read", &buf) + w := httptest.NewRecorder() + handler.readSearchLogFile(w, req, httprouter.Params{}) + if w.Code != 200 { + t.Fatalf("expected 200, got %d: %s", w.Code, w.Body.String()) + } + body := util.MapStr{} + if err := util.FromJSONBytes(w.Body.Bytes(), &body); err != nil { + t.Fatalf("failed to decode response: %v", err) + } + if body["success"] != true { + t.Errorf("expected success=true: %s", w.Body.String()) + } +} + +func quoteJSON(dir string) string { + return `"` + strings.ReplaceAll(strings.ReplaceAll(dir, `\`, `\\`), `"`, `\"`) + `"` } diff --git a/plugin/api/log_whitelist.go b/plugin/api/log_whitelist.go new file mode 100644 index 0000000..e66c50e --- /dev/null +++ b/plugin/api/log_whitelist.go @@ -0,0 +1,295 @@ +/* Copyright © INFINI Ltd. All rights reserved. + * Web: https://infinilabs.com + * Email: hello#infini.ltd */ + +package api + +import ( + "fmt" + "path/filepath" + "regexp" + "strings" + "sync" + "sync/atomic" + "time" + + "infini.sh/agent/lib/process" + readguard "infini.sh/agent/lib/util" + "infini.sh/framework/core/elastic" + "infini.sh/framework/core/env" + log "infini.sh/framework/core/log" + "infini.sh/framework/core/util" +) + +// Whitelist for the Elasticsearch/Easysearch log viewing API +// (/elasticsearch/logs/_list and /elasticsearch/logs/_read). These +// endpoints receive logs_path from the caller, which used to allow +// reading any directory visible to the agent process. Reads are now +// confined to a whitelist resolved from three sources: +// +// 1. elasticsearch_logs.allowed_paths in the agent config — the escape +// hatch for layouts where discovery cannot see the log directory +// 2. log directories the local search nodes report themselves +// (settings path.logs, plus path.home/logs), discovered with the +// same process scan the console gets its paths from +// 3. directories derived from the search processes' command lines +// (-Des.path.logs, path.home/logs, and the -Xlog gc file location), +// mirroring how the console derives the paths it sends back +// +// System paths (readguard.IsSystemReadPath) are never readable, even when +// whitelisted. The whitelist is cached for esLogDirsCacheTTL; a stale +// whitelist keeps serving while a refresh runs in the background. + +const esLogDirsCacheTTL = time.Minute + +// ESLogsConfig is the optional elasticsearch_logs config section. +type ESLogsConfig struct { + AllowedPaths []string `config:"allowed_paths" json:"allowed_paths"` +} + +var ( + esLogWhitelistMu sync.Mutex + esLogWhitelistGuard *readguard.ReadGuard + esLogWhitelistFetched time.Time + esLogWhitelistRefreshing atomic.Bool + + // esLogWhitelistLoader resolves the allowed roots; injectable in tests. + esLogWhitelistLoader = defaultESLogWhitelist +) + +// esLogsReadGuard returns the cached whitelist guard. When the cache is +// stale, the first caller refreshes it inline (single-flight via CAS) +// while concurrent callers keep being served the cached guard, so only +// one request per TTL pays the discovery cost; the first caller (no +// cache yet) also builds synchronously. When no whitelist can be +// established at all, access is denied (secure default). +func esLogsReadGuard() (*readguard.ReadGuard, error) { + esLogWhitelistMu.Lock() + guard := esLogWhitelistGuard + if guard != nil { + stale := time.Since(esLogWhitelistFetched) >= esLogDirsCacheTTL + esLogWhitelistMu.Unlock() + if !stale { + return guard, nil + } + if esLogWhitelistRefreshing.CompareAndSwap(false, true) { + defer esLogWhitelistRefreshing.Store(false) + refreshed, err := refreshESLogWhitelist() + if err != nil { + log.Warnf("failed to refresh elasticsearch logs whitelist, keeping the previous one: %v", err) + } + if refreshed != nil { + return refreshed, nil + } + } + return guard, nil + } + esLogWhitelistMu.Unlock() + return refreshESLogWhitelist() +} + +func refreshESLogWhitelist() (*readguard.ReadGuard, error) { + guard, err := buildESLogsReadGuard() + esLogWhitelistMu.Lock() + defer esLogWhitelistMu.Unlock() + if err != nil { + if esLogWhitelistGuard != nil { + return esLogWhitelistGuard, nil + } + return nil, err + } + esLogWhitelistGuard = guard + esLogWhitelistFetched = time.Now() + return guard, nil +} + +func resetESLogWhitelistCache() { + esLogWhitelistMu.Lock() + defer esLogWhitelistMu.Unlock() + esLogWhitelistGuard = nil +} + +func buildESLogsReadGuard() (*readguard.ReadGuard, error) { + roots, err := esLogWhitelistLoader() + if err != nil { + return nil, err + } + if len(roots) == 0 { + return nil, fmt.Errorf("no allowed elasticsearch log directories: configure elasticsearch_logs.allowed_paths or make sure the local search node is discoverable") + } + // validate roots one by one: a stale or misconfigured entry must not + // take the whole whitelist down + var valid []string + seen := map[string]bool{} + for _, root := range roots { + if seen[root] { + continue + } + seen[root] = true + if _, err := readguard.NewReadGuard(root); err != nil { + log.Warnf("ignoring invalid elasticsearch logs path [%s]: %v", root, err) + continue + } + valid = append(valid, root) + } + if len(valid) == 0 { + return nil, fmt.Errorf("no valid elasticsearch log directories in %v", roots) + } + return readguard.NewReadGuard(valid...) +} + +// defaultESLogWhitelist combines the static config section with the log +// directories discovered from the local search nodes and their command +// lines. +func defaultESLogWhitelist() ([]string, error) { + return append(append(loadStaticESLogPaths(), discoverESLogDirs()...), cmdlineESLogDirs()...), nil +} + +// loadStaticESLogPaths reads elasticsearch_logs.allowed_paths from the +// app config file (main file + config dir), falling back to the global +// env config like the discovery endpoints do. +func loadStaticESLogPaths() []string { + cfg := ESLogsConfig{} + var err error + appCfg, cfgErr := getAppConfig() + if cfgErr != nil { + _, err = env.ParseConfig("elasticsearch_logs", &cfg) + } else { + _, err = env.ParseConfigSection(appCfg, "elasticsearch_logs", &cfg) + } + if err != nil { + log.Debugf("no elasticsearch_logs config: %v", err) + return nil + } + return cfg.AllowedPaths +} + +// discoverESLogDirs runs the same local process scan the console-facing +// discovery endpoint uses, and extracts each node's reported log dir. +func discoverESLogDirs() []string { + result, err := process.DiscoverESNode(nil) + if err != nil { + log.Warnf("failed to discover local search nodes for logs whitelist: %v", err) + return nil + } + var dirs []string + for _, node := range result.Nodes { + dirs = append(dirs, nodeLogDirs(node.NodeInfo)...) + } + return dirs +} + +// nodeLogDirs extracts settings path.logs (string or array) plus +// path.home/logs, so the whitelist covers every location the console can +// derive from the node's settings. +func nodeLogDirs(info *elastic.NodesInfo) []string { + if info == nil || len(info.Settings) == 0 { + return nil + } + settings := util.MapStr(info.Settings) + var dirs []string + if v, err := settings.GetValue("path.logs"); err == nil { + dirs = append(dirs, pathList(v)...) + } + if v, err := settings.GetValue("path.home"); err == nil { + if home, err := util.ExtractString(v); err == nil && home != "" { + dirs = append(dirs, filepath.Join(home, "logs")) + } + } + return dirs +} + +// cmdlineESLogDirs mirrors the console's deriveLogsPathsFromCmdline so +// that whatever directory the console computes from a process command +// line (-Des.path.logs, path.home/logs, -Xlog gc file dir) is in the +// whitelist before it can be requested back. +func cmdlineESLogDirs() []string { + procs, err := process.DiscoverESProcessors(process.ElasticFilter) + if err != nil { + log.Warnf("failed to scan search processes for logs whitelist: %v", err) + return nil + } + var dirs []string + for _, p := range procs { + dirs = append(dirs, cmdlineLogDirs(p.Cmdline)...) + } + return dirs +} + +var ( + cmdlinePathHomeRe = regexp.MustCompile(`(?:^|\s)-D(?:es|opensearch)\.path\.home=([^\s]+)`) + cmdlinePathLogsRe = regexp.MustCompile(`(?:^|\s)-D(?:es|opensearch)\.path\.logs=([^\s]+)`) + cmdlineGCFileRe = regexp.MustCompile(`(?:^|\s)-Xlog:[^\s]*?file=([^\s]+)`) +) + +func cmdlineLogDirs(cmdline string) []string { + pathHome := cmdlineValue(cmdlinePathHomeRe, cmdline) + var dirs []string + if v := cmdlineValue(cmdlinePathLogsRe, cmdline); v != "" { + if r := resolveCmdlinePath(v, pathHome); r != "" { + dirs = append(dirs, r) + } + } else if pathHome != "" { + dirs = append(dirs, filepath.Join(pathHome, "logs")) + } + if v := trimGCLogFileValue(cmdlineValue(cmdlineGCFileRe, cmdline)); v != "" { + if r := resolveCmdlinePath(v, pathHome); r != "" { + dirs = append(dirs, filepath.Dir(r)) + } + } + return dirs +} + +func cmdlineValue(re *regexp.Regexp, cmdline string) string { + matches := re.FindStringSubmatch(cmdline) + if len(matches) > 1 { + return strings.Trim(strings.TrimSpace(matches[1]), `"'`) + } + return "" +} + +// trimGCLogFileValue drops the rotation tags after the file name, e.g. +// /var/log/gc.log:uptime,tags -> /var/log/gc.log (drive letters kept). +func trimGCLogFileValue(value string) string { + searchFrom := 0 + if len(value) > 1 && value[1] == ':' { + searchFrom = 2 + } + if idx := strings.Index(value[searchFrom:], ":"); idx >= 0 { + value = value[:searchFrom+idx] + } + return value +} + +func resolveCmdlinePath(value, base string) string { + if value == "" { + return "" + } + if !filepath.IsAbs(value) { + if base == "" { + return "" + } + value = filepath.Join(base, value) + } + return filepath.Clean(value) +} + +func pathList(raw interface{}) []string { + switch v := raw.(type) { + case string: + if v != "" { + return []string{v} + } + case []string: + return v + case []interface{}: + items := make([]string, 0, len(v)) + for _, item := range v { + if s := util.ToString(item); s != "" { + items = append(items, s) + } + } + return items + } + return nil +} diff --git a/plugin/logs/file_detect.go b/plugin/logs/file_detect.go index b694a0d..ec7a659 100644 --- a/plugin/logs/file_detect.go +++ b/plugin/logs/file_detect.go @@ -31,21 +31,30 @@ type FSEvent struct { State FileState } +type FileDetector struct { + root string + patterns []*Pattern + prev map[string]os.FileInfo + events chan FSEvent + + // known tracks the files seen in the last walk; files that vanish + // become stale candidates so a re-created path carrying the same + // file identity (rename/move) can inherit their offset instead of + // being re-read from the beginning. + known map[string]bool + stale map[string]FileState +} + func NewFileDetector(rootPath string, patterns []*Pattern) *FileDetector { return &FileDetector{ root: rootPath, patterns: patterns, events: make(chan FSEvent), + known: map[string]bool{}, + stale: map[string]FileState{}, } } -type FileDetector struct { - root string - patterns []*Pattern - prev map[string]os.FileInfo - events chan FSEvent -} - func (w *FileDetector) Detect(ctx context.Context) { defer func() { w.events <- doneEvent() @@ -54,6 +63,7 @@ func (w *FileDetector) Detect(ctx context.Context) { if len(w.patterns) == 0 { return } + walked := map[string]bool{} err := filepath.Walk(w.root, func(path string, info os.FileInfo, err error) error { if ctx.Err() != nil { return ctx.Err() @@ -70,6 +80,8 @@ func (w *FileDetector) Detect(ctx context.Context) { continue } w.judgeEvent(ctx, path, info, pattern) + walked[path] = true + w.known[path] = true break } return nil @@ -77,6 +89,16 @@ func (w *FileDetector) Detect(ctx context.Context) { if err != nil { log.Errorf("failed to walk logs under [%s], err: %v", w.root, err) } + + // refresh rename candidates: known files missing from this walk + for path := range w.known { + if !walked[path] { + if state, err := GetFileState(path); err == nil && state != (FileState{}) { + w.stale[path] = state + } + delete(w.known, path) + } + } } func (w *FileDetector) judgeEvent(ctx context.Context, path string, info os.FileInfo, pattern *Pattern) { @@ -96,6 +118,16 @@ func (w *FileDetector) judgeEvent(ctx context.Context, path string, info os.File preState, err := GetFileState(path) isSameFile := w.IsSameFile(preState, info, path) if err != nil || preState == (FileState{}) || !isSameFile { + // rename/move: a vanished file with the same identity hands over + // its offset so the content is not re-ingested + for stalePath, staleState := range w.stale { + if stalePath != path && w.IsSameFile(staleState, info, path) { + preState = staleState + delete(w.stale, stalePath) + log.Debugf("file moved: %s -> %s, inheriting offset %d", stalePath, path, preState.Offset) + break + } + } select { case <-ctx.Done(): return diff --git a/plugin/logs/logs.go b/plugin/logs/logs.go index bacdbe9..8c2a5c6 100644 --- a/plugin/logs/logs.go +++ b/plugin/logs/logs.go @@ -20,7 +20,6 @@ import ( "infini.sh/framework/core/global" log "infini.sh/framework/core/log" "infini.sh/framework/core/pipeline" - "infini.sh/framework/core/queue" "infini.sh/framework/core/task" "infini.sh/framework/core/util" ) @@ -29,6 +28,7 @@ type LogsProcessor struct { cfg Config watcher *FileDetector agentMeta *event2.AgentMeta + emit *emitter lock sync.RWMutex } @@ -54,10 +54,42 @@ type Pattern struct { } type Config struct { - QueueName string `config:"queue_name"` + QueueName string `config:"queue_name"` + // QueueType selects the queue backend (empty = default disk). Set + // to "kafka" to write logs straight to the Kafka bus (brokers come + // from the instance-level kafka_queue section), consumed by Gateway. + QueueType string `config:"queue_type"` LogsPath string `config:"logs_path"` Metadata util.MapStr `config:"metadata"` Patterns []*Pattern `config:"patterns"` + + // ScanInterval controls how often the logs path is re-scanned while + // the processor runs (default 10s); lower it for closer-to-live + // tailing. Empty means a single scan per pipeline run (legacy + // behavior). + ScanInterval string `config:"scan_interval"` + + // ShipDirect bypasses the local queue: envelopes are shipped + // straight to the configured shipper (default: OTLP/gRPC to the + // gateway tier). The file itself plus offset checkpoints provide + // the durability -- offsets only advance after successful delivery + // -- so the double disk I/O of a local queue copy is avoided. + ShipDirect bool `config:"ship_direct"` + + // Shipper names the direct-ship transport (default "otlp"). + Shipper string `config:"shipper"` + + // ShipBatchSize flushes the in-flight batch at this many events + // (default 500). Memory bound of ship mode. + ShipBatchSize int `config:"ship_batch_size"` + + // ShipFlushInterval flushes the in-flight batch at least this often + // while a file is being read (default 1s). + ShipFlushInterval string `config:"ship_flush_interval"` + + // ShipConfig is passed through to the shipper factory (for "otlp": + // the same keys as the otlp_export processor). + ShipConfig map[string]interface{} `config:"ship_config"` } func init() { @@ -90,9 +122,14 @@ func NewFromConfig(cfg Config) (pipeline.Processor, error) { } patterns = append(patterns, pattern) } + emit, err := newEmitter(cfg) + if err != nil { + return nil, err + } p := &LogsProcessor{ cfg: cfg, watcher: NewFileDetector(cfg.LogsPath, cfg.Patterns), + emit: emit, } return p, nil @@ -113,17 +150,48 @@ func (p *LogsProcessor) Name() string { } func (p *LogsProcessor) Process(c *pipeline.Context) error { - task.RunWithinGroup(name, func(ctx context.Context) error { - p.watcher.Detect(ctx) - return nil - }) - var fsEvent FSEvent + interval := time.Duration(0) + if p.cfg.ScanInterval != "" { + d, err := time.ParseDuration(p.cfg.ScanInterval) + if err != nil || d <= 0 { + return fmt.Errorf("%s: invalid scan_interval %q", name, p.cfg.ScanInterval) + } + interval = d + } + + // derive from the pipeline context so an in-flight walk aborts when + // the pipeline stops; no extra goroutine is spawned (the embedded + // stdlib cancelCtx is linked via the parentCancelCtx fast path) + scanCtx, cancel := context.WithCancel(c) + defer cancel() + + first := true for !c.IsCanceled() { - fsEvent = p.watcher.Event() - if fsEvent.Op == OpDone { + task.RunWithinGroup(name, func(ctx context.Context) error { + // the detector signals completion with a done event; use a + // context tied to this scan run + p.watcher.Detect(scanCtx) return nil + }) + if first { + first = false + } + // drain all pending events of this scan before the next walk + for !c.IsCanceled() { + fsEvent := p.watcher.Event() + if fsEvent.Op == OpDone { + break + } + p.onFSEvent(fsEvent, c) + } + if interval <= 0 { + return nil // legacy single-scan behavior + } + select { + case <-c.Done(): + return nil + case <-time.After(interval): } - p.onFSEvent(fsEvent, c) } return nil } @@ -163,6 +231,7 @@ func (p *LogsProcessor) ReadLogs(event FSEvent, c *pipeline.Context) { func (p *LogsProcessor) ReadJsonLogs(event FSEvent, c *pipeline.Context) { log.Debugf("reading json logs from [%s], offset: [%d]", event.Path, event.State.Offset) offset := event.State.Offset + p.emit.beginFile(offset) h, err := harvester.NewHarvester(event.Path, offset) if err != nil { log.Errorf("failed to initialize harvester, err: %v", err) @@ -174,6 +243,10 @@ func (p *LogsProcessor) ReadJsonLogs(event FSEvent, c *pipeline.Context) { return } for !c.IsCanceled() { + if err := p.emit.maybeFlush(); err != nil { + log.Errorf("failed to flush batch for file [%s], err: %v", event.Path, err) + break + } msg, err := r.Next() if err == io.EOF { break @@ -191,22 +264,11 @@ func (p *LogsProcessor) ReadJsonLogs(event FSEvent, c *pipeline.Context) { continue } logContent, timestamp := processJSON(event.Pattern, logContent) - p.Save(event, logContent, timestamp) - } - sysInfo, err := LoadFileID(event.Info, event.Path) - if err != nil { - log.Errorf("failed to get file info, err: %v", err) - return - } - event.State = FileState{ - Name: event.Info.Name(), - Size: event.Info.Size(), - ModTime: event.Info.ModTime(), - Path: event.Path, - Offset: offset, - Sys: sysInfo, + if !p.emitEvent(event, logContent, timestamp, offset) { + break + } } - SaveFileState(event.Path, event.State) + p.finishFile(event) } func (p *LogsProcessor) ReadPlainTextLogs(event FSEvent, c *pipeline.Context) { @@ -222,8 +284,13 @@ func (p *LogsProcessor) ReadPlainTextLogs(event FSEvent, c *pipeline.Context) { return } offset := event.State.Offset + p.emit.beginFile(offset) var logMessage string for !c.IsCanceled() { + if err := p.emit.maybeFlush(); err != nil { + log.Errorf("failed to flush batch for file [%s], err: %v", event.Path, err) + break + } msg, err := r.Next() if err == io.EOF { break @@ -239,23 +306,11 @@ func (p *LogsProcessor) ReadPlainTextLogs(event FSEvent, c *pipeline.Context) { logMessage = util.UnsafeBytesToString(msg.Content) logContent, timestamp := processText(event.Pattern, logMessage) logContent["message"] = logMessage - p.Save(event, logContent, timestamp) - } - - sysInfo, err := LoadFileID(event.Info, event.Path) - if err != nil { - log.Errorf("failed to get file info, err: %v", err) - return - } - event.State = FileState{ - Name: event.Info.Name(), - Size: event.Info.Size(), - ModTime: event.Info.ModTime(), - Path: event.Path, - Offset: offset, - Sys: sysInfo, + if !p.emitEvent(event, logContent, timestamp, offset) { + break + } } - SaveFileState(event.Path, event.State) + p.finishFile(event) } func (p *LogsProcessor) ReadMultilineLogs(event FSEvent, c *pipeline.Context) { @@ -271,8 +326,13 @@ func (p *LogsProcessor) ReadMultilineLogs(event FSEvent, c *pipeline.Context) { return } offset := event.State.Offset + p.emit.beginFile(offset) var logMessage string for !c.IsCanceled() { + if err := p.emit.maybeFlush(); err != nil { + log.Errorf("failed to flush batch for file [%s], err: %v", event.Path, err) + break + } msg, err := r.Next() if err == io.EOF { break @@ -288,26 +348,16 @@ func (p *LogsProcessor) ReadMultilineLogs(event FSEvent, c *pipeline.Context) { logMessage = util.UnsafeBytesToString(msg.Content) logContent, timestamp := processText(event.Pattern, logMessage) logContent["message"] = logMessage - p.Save(event, logContent, timestamp) - } - - sysInfo, err := LoadFileID(event.Info, event.Path) - if err != nil { - log.Errorf("failed to get file info, err: %v", err) - return - } - event.State = FileState{ - Name: event.Info.Name(), - Size: event.Info.Size(), - ModTime: event.Info.ModTime(), - Path: event.Path, - Offset: offset, - Sys: sysInfo, + if !p.emitEvent(event, logContent, timestamp, offset) { + break + } } - SaveFileState(event.Path, event.State) + p.finishFile(event) } -func (p *LogsProcessor) Save(event FSEvent, logContent util.MapStr, timestamp string) { +// buildEnvelope renders one collected event as the LogEvent envelope +// JSON (the format shared with the queue boundary and the gateway). +func (p *LogsProcessor) buildEnvelope(event FSEvent, logContent util.MapStr, timestamp string) []byte { logEvent := LogEvent{ AgentMeta: p.GetAgentMeta(), Fields: logContent, @@ -331,7 +381,44 @@ func (p *LogsProcessor) Save(event FSEvent, logContent util.MapStr, timestamp st } else { logEvent.Timestamp = time.Now().Format(time.RFC3339) } - queue.Push(queue.GetOrInitConfig(logEvent.AgentMeta.LoggingQueueName), util.MustToJSONBytes(logEvent)) + return util.MustToJSONBytes(logEvent) +} + +// emitEvent delivers one event via the emitter (queue or direct +// shipper). It returns false when delivery failed and the read loop +// must stop; the file state stays at the committed offset so the event +// is re-delivered on the next scan (at-least-once). +func (p *LogsProcessor) emitEvent(event FSEvent, logContent util.MapStr, timestamp string, endOffset int64) bool { + data := p.buildEnvelope(event, logContent, timestamp) + if err := p.emit.emit(data, endOffset); err != nil { + log.Errorf("failed to deliver log event from file [%s] at offset %d, will retry on next scan: %v", + event.Path, endOffset, err) + return false + } + return true +} + +// finishFile flushes the in-flight ship batch (if any) and persists the +// file state at the committed offset. +func (p *LogsProcessor) finishFile(event FSEvent) { + if err := p.emit.flush(); err != nil { + log.Errorf("failed to flush batch for file [%s], state stays at offset %d: %v", + event.Path, p.emit.committed, err) + } + sysInfo, err := LoadFileID(event.Info, event.Path) + if err != nil { + log.Errorf("failed to get file info, err: %v", err) + return + } + event.State = FileState{ + Name: event.Info.Name(), + Size: event.Info.Size(), + ModTime: event.Info.ModTime(), + Path: event.Path, + Offset: p.emit.committed, + Sys: sysInfo, + } + SaveFileState(event.Path, event.State) } func (p *LogsProcessor) GetAgentMeta() *event2.AgentMeta { diff --git a/plugin/logs/ship.go b/plugin/logs/ship.go new file mode 100644 index 0000000..d4fc208 --- /dev/null +++ b/plugin/logs/ship.go @@ -0,0 +1,149 @@ +/* Copyright © INFINI Ltd. All rights reserved. + * Web: https://infinilabs.com + * Email: hello#infini.ltd */ + +package logs + +import ( + "fmt" + "time" + + "infini.sh/framework/core/queue" + "infini.sh/framework/core/shipper" +) + +const ( + defaultShipBatchSize = 500 + defaultShipFlushInterval = time.Second +) + +// emitter routes log-event envelopes to either the local queue (durable +// buffering; the default) or a direct shipper (ship_direct mode). +// +// In ship mode the file itself plus offset checkpoints provide the +// durability: envelopes are batched in memory (bounded by batchSize), +// and the committed offset only advances after a batch is delivered. +// On delivery failure the pending batch is dropped and the file is +// re-read from the committed offset on the next scan -- the local disk +// queue is not needed. +// +// Harvesting is strictly sequential (one file at a time in the Process +// loop), so the emitter keeps a single in-flight batch without locks. +type emitter struct { + shipMode bool + queueName string + queueType string // empty = default disk; "kafka" = write directly to the Kafka bus + batchSize int + flushEvery time.Duration + + shipper shipper.Shipper + ticker *time.Ticker + + batch [][]byte // in-flight envelopes (ship mode) + offsets []int64 // end-offset of each envelope in the batch + committed int64 // end-offset of the last delivered event +} + +// newEmitter builds the emitter described by cfg; queue mode needs no +// external resources, ship mode resolves the named shipper. +func newEmitter(cfg Config) (*emitter, error) { + e := &emitter{ + shipMode: cfg.ShipDirect, + queueName: cfg.QueueName, + queueType: cfg.QueueType, + batchSize: cfg.ShipBatchSize, + } + if !e.shipMode { + return e, nil + } + + if e.batchSize <= 0 { + e.batchSize = defaultShipBatchSize + } + if cfg.ShipFlushInterval != "" { + d, err := time.ParseDuration(cfg.ShipFlushInterval) + if err != nil || d <= 0 { + return nil, fmt.Errorf("%s: invalid ship_flush_interval %q", name, cfg.ShipFlushInterval) + } + e.flushEvery = d + } else { + e.flushEvery = defaultShipFlushInterval + } + + shipperName := cfg.Shipper + if shipperName == "" { + shipperName = "otlp" + } + s, err := shipper.Get(shipperName, cfg.ShipConfig) + if err != nil { + return nil, fmt.Errorf("%s: ship_direct requires the %q shipper: %v", name, shipperName, err) + } + e.shipper = s + e.ticker = time.NewTicker(e.flushEvery) + return e, nil +} + +// beginFile resets the per-file batch state; committed starts at the +// file's current offset. +func (e *emitter) beginFile(startOffset int64) { + e.batch = e.batch[:0] + e.offsets = e.offsets[:0] + e.committed = startOffset +} + +// emit delivers one envelope. It returns an error when delivery failed; +// the caller must stop reading the file and persist the committed +// offset so the event is re-delivered on the next scan (at-least-once). +func (e *emitter) emit(data []byte, endOffset int64) error { + if !e.shipMode { + // EnsureTypedConfig: force-register the backend when queue_type + // names one explicitly (kafka bus mode), otherwise create the + // default disk queue on demand. A failed Push returns the error + // so the caller stops reading and keeps its offset — natural + // backpressure down to the harvester. + qcfg := queue.EnsureTypedConfig(e.queueType, e.queueName) + if err := queue.Push(qcfg, data); err != nil { + return err + } + e.committed = endOffset + return nil + } + + e.batch = append(e.batch, data) + e.offsets = append(e.offsets, endOffset) + if len(e.batch) >= e.batchSize { + return e.flush() + } + return nil +} + +// maybeFlush ships the in-flight batch when the flush interval has +// elapsed (non-blocking; call between reads). +func (e *emitter) maybeFlush() error { + if !e.shipMode || e.ticker == nil || len(e.batch) == 0 { + return nil + } + select { + case <-e.ticker.C: + return e.flush() + default: + return nil + } +} + +// flush ships the in-flight batch; on success the committed offset +// advances to the last delivered event's end offset. On failure the +// batch is kept so the caller can see the error; the events will be +// re-read from the file (the durable buffer) on the next scan. +func (e *emitter) flush() error { + if !e.shipMode || len(e.batch) == 0 { + return nil + } + if err := e.shipper.Ship(e.batch); err != nil { + return err + } + e.committed = e.offsets[len(e.offsets)-1] + e.batch = e.batch[:0] + e.offsets = e.offsets[:0] + return nil +} diff --git a/plugin/logs/ship_test.go b/plugin/logs/ship_test.go new file mode 100644 index 0000000..cc13ced --- /dev/null +++ b/plugin/logs/ship_test.go @@ -0,0 +1,176 @@ +/* Copyright © INFINI Ltd. All rights reserved. + * Web: https://infinilabs.com + * Email: hello#infini.ltd */ + +package logs + +import ( + "fmt" + "sync" + "testing" + + "infini.sh/framework/core/shipper" +) + +// mockShipper records the batches it was asked to deliver; failN makes +// the first N Ship calls fail. +type mockShipper struct { + mu sync.Mutex + batches []int // size of each delivered batch + failN int + calls int + closeCnt int +} + +func (m *mockShipper) Ship(batch [][]byte) error { + m.mu.Lock() + defer m.mu.Unlock() + m.calls++ + if m.calls <= m.failN { + return fmt.Errorf("mock ship failure %d", m.calls) + } + m.batches = append(m.batches, len(batch)) + return nil +} + +func (m *mockShipper) Close() error { + m.mu.Lock() + defer m.mu.Unlock() + m.closeCnt++ + return nil +} + +var ( + registerOnce sync.Once + currentMock *mockShipper +) + +// useMockShipper registers the "mock_ship_test" factory once (Go test +// binaries run tests in one process; duplicate registry entries panic) +// and swaps in a fresh mock for the next emitter construction. +func useMockShipper(failN int) *mockShipper { + registerOnce.Do(func() { + shipper.Register("mock_ship_test", func(cfg map[string]interface{}) (shipper.Shipper, error) { + return currentMock, nil + }) + }) + m := &mockShipper{failN: failN} + currentMock = m + return m +} + +func shipModeConfig(batchSize int) Config { + return Config{ + QueueName: "logs", + ShipDirect: true, + Shipper: "mock_ship_test", + ShipBatchSize: batchSize, + } +} + +func data(i int) []byte { return []byte(fmt.Sprintf("event-%d", i)) } + +// TestEmitterFlushOnBatchSize verifies the in-flight batch is shipped +// as soon as it reaches batchSize and the committed offset advances to +// the last delivered event. +func TestEmitterFlushOnBatchSize(t *testing.T) { + m := useMockShipper(0) + e, err := newEmitter(shipModeConfig(2)) + if err != nil { + t.Fatalf("newEmitter: %v", err) + } + + e.beginFile(100) + if err := e.emit(data(1), 110); err != nil { + t.Fatalf("emit 1: %v", err) + } + if err := e.emit(data(2), 120); err != nil { + t.Fatalf("emit 2: %v", err) // batch full -> auto flush + } + if got := e.committed; got != 120 { + t.Fatalf("committed = %d, want 120", got) + } + if len(m.batches) != 1 || m.batches[0] != 2 { + t.Fatalf("batches = %v, want [2]", m.batches) + } + + // remainder ships on the final flush + if err := e.emit(data(3), 130); err != nil { + t.Fatalf("emit 3: %v", err) + } + if err := e.flush(); err != nil { + t.Fatalf("flush: %v", err) + } + if got := e.committed; got != 130 { + t.Fatalf("committed after final flush = %d, want 130", got) + } + if len(m.batches) != 2 || m.batches[1] != 1 { + t.Fatalf("batches = %v, want [2 1]", m.batches) + } +} + +// TestEmitterFailureKeepsOffset verifies that when delivery fails the +// committed offset does NOT advance -- the caller re-reads from the +// committed offset on the next scan (the file is the durable buffer). +func TestEmitterFailureKeepsOffset(t *testing.T) { + useMockShipper(1) // first Ship call fails + e, err := newEmitter(shipModeConfig(1)) + if err != nil { + t.Fatalf("newEmitter: %v", err) + } + + e.beginFile(50) + if err := e.emit(data(1), 60); err == nil { + t.Fatal("expected delivery error, got nil") + } + if got := e.committed; got != 50 { + t.Fatalf("committed = %d, want 50 (unchanged after failure)", got) + } +} + +// TestEmitterQueueModeRequiresNoShipper verifies queue mode builds +// without any registered shipper (no external dependencies). +func TestEmitterQueueModeRequiresNoShipper(t *testing.T) { + e, err := newEmitter(Config{QueueName: "logs"}) + if err != nil { + t.Fatalf("newEmitter (queue mode): %v", err) + } + if e.shipMode { + t.Fatal("default mode must be queue mode") + } +} + +// TestEmitterBeginFileResetsBatch verifies per-file batch state resets +// and the pending (undelivered) batch of the previous file is dropped: +// those bytes will be re-read from the file. +func TestEmitterBeginFileResetsBatch(t *testing.T) { + useMockShipper(0) + e, err := newEmitter(shipModeConfig(100)) // never auto-flushes + if err != nil { + t.Fatalf("newEmitter: %v", err) + } + + e.beginFile(0) + _ = e.emit(data(1), 10) // stays in-flight + e.beginFile(500) + if len(e.batch) != 0 { + t.Fatalf("batch not reset: %d events", len(e.batch)) + } + if got := e.committed; got != 500 { + t.Fatalf("committed = %d, want 500", got) + } +} + +// TestEmitterInvalidFlushInterval verifies a bad ship_flush_interval is +// rejected at construction time. +func TestEmitterInvalidFlushInterval(t *testing.T) { + useMockShipper(0) + _, err := newEmitter(Config{ + ShipDirect: true, + Shipper: "mock_ship_test", + ShipFlushInterval: "nope", + }) + if err == nil { + t.Fatal("expected error for invalid ship_flush_interval, got nil") + } +} diff --git a/plugin/plugin.go b/plugin/plugin.go new file mode 100644 index 0000000..8b9cb98 --- /dev/null +++ b/plugin/plugin.go @@ -0,0 +1,18 @@ +/* Copyright © INFINI Ltd. All rights reserved. + * Web: https://infinilabs.com + * Email: hello#infini.ltd */ + +// Package plugin aggregates the agent-side plugin packages so that their +// init() side effects (pipeline processor registrations, module hooks) run +// when the agent boots. main.go blank-imports this package. +package plugin + +import ( + _ "infini.sh/agent/plugin/elastic/logging" + _ "infini.sh/agent/plugin/elastic/metric" + _ "infini.sh/agent/plugin/logs" + + // kafka queue backend: enables routing the "logs" queue to Kafka + // purely via configuration (kafka_queue.default: true) + _ "infini.sh/framework/plugins/queue/kafka_queue" +)