@cryptotaxi247 / netdata-1 / commits / 7d312bb23

chore(otel/journaldexporter): improve remote tests (#20143)

Ilya Mashchenko committed Apr 18, 2025 at 16:01 UTC 7d312bb23f6eda2aa400559e08771e8844259a5b
5 files changed +420 -473
src/go/otel-collector/builder-config.yaml
+5 -5
@@ -8,13 +8,13 @@ dist:
8
9
10 receivers:
11 - - gomod: go.opentelemetry.io/collector/receiver/otlpreceiver v0.123.0
12 - - gomod: github.com/open-telemetry/opentelemetry-collector-contrib/receiver/filelogreceiver v0.123.0
11 + - gomod: go.opentelemetry.io/collector/receiver/otlpreceiver v0.124.0
12 + - gomod: github.com/open-telemetry/opentelemetry-collector-contrib/receiver/filelogreceiver v0.124.0
13
14 exporters:
15 - - gomod: go.opentelemetry.io/collector/exporter/debugexporter v0.123.0
16 - - gomod: github.com/open-telemetry/opentelemetry-collector-contrib/exporter/fileexporter v0.123.0
17 - - gomod: go.opentelemetry.io/collector/exporter/otlpexporter v0.123.0
15 + - gomod: go.opentelemetry.io/collector/exporter/debugexporter v0.124.0
16 + - gomod: github.com/open-telemetry/opentelemetry-collector-contrib/exporter/fileexporter v0.124.0
17 + - gomod: go.opentelemetry.io/collector/exporter/otlpexporter v0.124.0
18 - gomod: github.com/netdata/netdata/otel-collector/exporter/journaldexporter v0.0.0
19
20 replaces:
src/go/otel-collector/exporter/journaldexporter/go.mod
+30 -30
@@ -3,15 +3,16 @@ module github.com/netdata/netdata/otel-collector/exporter/journaldexporter
3 go 1.23.4
4
5 require (
6 + github.com/google/uuid v1.6.0
7 github.com/stretchr/testify v1.10.0
7 - go.opentelemetry.io/collector/component v1.29.0
8 - go.opentelemetry.io/collector/component/componenttest v0.123.0
9 - go.opentelemetry.io/collector/confmap v1.29.0
10 - go.opentelemetry.io/collector/consumer v1.29.0
11 - go.opentelemetry.io/collector/exporter v0.123.0
12 - go.opentelemetry.io/collector/exporter/exportertest v0.123.0
13 - go.opentelemetry.io/collector/exporter/xexporter v0.123.0
14 - go.opentelemetry.io/collector/pdata v1.29.0
8 + go.opentelemetry.io/collector/component v1.30.0
9 + go.opentelemetry.io/collector/component/componenttest v0.124.0
10 + go.opentelemetry.io/collector/confmap v1.30.0
11 + go.opentelemetry.io/collector/consumer v1.30.0
12 + go.opentelemetry.io/collector/exporter v0.124.0
13 + go.opentelemetry.io/collector/exporter/exportertest v0.124.0
14 + go.opentelemetry.io/collector/exporter/xexporter v0.124.0
15 + go.opentelemetry.io/collector/pdata v1.30.0
16 go.uber.org/goleak v1.3.0
17 go.uber.org/zap v1.27.0
18 )
@@ -23,31 +24,30 @@ require (
24 github.com/go-logr/stdr v1.2.2 // indirect
25 github.com/go-viper/mapstructure/v2 v2.2.1 // indirect
26 github.com/gogo/protobuf v1.3.2 // indirect
26 - github.com/google/uuid v1.6.0 // indirect
27 github.com/hashicorp/go-version v1.7.0 // indirect
28 github.com/json-iterator/go v1.1.12 // indirect
29 - github.com/knadh/koanf/maps v0.1.1 // indirect
30 - github.com/knadh/koanf/providers/confmap v0.1.0 // indirect
31 - github.com/knadh/koanf/v2 v2.1.2 // indirect
29 + github.com/knadh/koanf/maps v0.1.2 // indirect
30 + github.com/knadh/koanf/providers/confmap v1.0.0 // indirect
31 + github.com/knadh/koanf/v2 v2.2.0 // indirect
32 github.com/mitchellh/copystructure v1.2.0 // indirect
33 github.com/mitchellh/reflectwalk v1.0.2 // indirect
34 github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd // indirect
35 github.com/modern-go/reflect2 v1.0.2 // indirect
36 github.com/pmezard/go-difflib v1.0.0 // indirect
37 go.opentelemetry.io/auto/sdk v1.1.0 // indirect
38 - go.opentelemetry.io/collector/config/configretry v1.29.0 // indirect
39 - go.opentelemetry.io/collector/consumer/consumererror v0.123.0 // indirect
40 - go.opentelemetry.io/collector/consumer/consumertest v0.123.0 // indirect
41 - go.opentelemetry.io/collector/consumer/xconsumer v0.123.0 // indirect
42 - go.opentelemetry.io/collector/extension v1.29.0 // indirect
43 - go.opentelemetry.io/collector/extension/xextension v0.123.0 // indirect
44 - go.opentelemetry.io/collector/featuregate v1.29.0 // indirect
45 - go.opentelemetry.io/collector/internal/telemetry v0.123.0 // indirect
46 - go.opentelemetry.io/collector/pdata/pprofile v0.123.0 // indirect
47 - go.opentelemetry.io/collector/pipeline v0.123.0 // indirect
48 - go.opentelemetry.io/collector/receiver v1.29.0 // indirect
49 - go.opentelemetry.io/collector/receiver/receivertest v0.123.0 // indirect
50 - go.opentelemetry.io/collector/receiver/xreceiver v0.123.0 // indirect
38 + go.opentelemetry.io/collector/config/configretry v1.30.0 // indirect
39 + go.opentelemetry.io/collector/consumer/consumererror v0.124.0 // indirect
40 + go.opentelemetry.io/collector/consumer/consumertest v0.124.0 // indirect
41 + go.opentelemetry.io/collector/consumer/xconsumer v0.124.0 // indirect
42 + go.opentelemetry.io/collector/extension v1.30.0 // indirect
43 + go.opentelemetry.io/collector/extension/xextension v0.124.0 // indirect
44 + go.opentelemetry.io/collector/featuregate v1.30.0 // indirect
45 + go.opentelemetry.io/collector/internal/telemetry v0.124.0 // indirect
46 + go.opentelemetry.io/collector/pdata/pprofile v0.124.0 // indirect
47 + go.opentelemetry.io/collector/pipeline v0.124.0 // indirect
48 + go.opentelemetry.io/collector/receiver v1.30.0 // indirect
49 + go.opentelemetry.io/collector/receiver/receivertest v0.124.0 // indirect
50 + go.opentelemetry.io/collector/receiver/xreceiver v0.124.0 // indirect
51 go.opentelemetry.io/contrib/bridges/otelzap v0.10.0 // indirect
52 go.opentelemetry.io/otel v1.35.0 // indirect
53 go.opentelemetry.io/otel/log v0.11.0 // indirect
@@ -56,11 +56,11 @@ require (
56 go.opentelemetry.io/otel/sdk/metric v1.35.0 // indirect
57 go.opentelemetry.io/otel/trace v1.35.0 // indirect
58 go.uber.org/multierr v1.11.0 // indirect
59 - golang.org/x/net v0.37.0 // indirect
60 - golang.org/x/sys v0.31.0 // indirect
61 - golang.org/x/text v0.23.0 // indirect
62 - google.golang.org/genproto/googleapis/rpc v0.0.0-20250303144028-a0af3efb3deb // indirect
63 - google.golang.org/grpc v1.71.0 // indirect
59 + golang.org/x/net v0.39.0 // indirect
60 + golang.org/x/sys v0.32.0 // indirect
61 + golang.org/x/text v0.24.0 // indirect
62 + google.golang.org/genproto/googleapis/rpc v0.0.0-20250414145226-207652e42e2e // indirect
63 + google.golang.org/grpc v1.71.1 // indirect
64 google.golang.org/protobuf v1.36.6 // indirect
65 gopkg.in/yaml.v3 v3.0.1 // indirect
66 )
src/go/otel-collector/exporter/journaldexporter/go.sum
+62 -62
@@ -25,12 +25,12 @@ github.com/json-iterator/go v1.1.12 h1:PV8peI4a0ysnczrg+LtxykD8LfKY9ML6u2jnxaEnr
25 github.com/json-iterator/go v1.1.12/go.mod h1:e30LSqwooZae/UwlEbR2852Gd8hjQvJoHmT4TnhNGBo=
26 github.com/kisielk/errcheck v1.5.0/go.mod h1:pFxgyoBC7bSaBwPgfKdkLd5X25qrDl4LWUI2bnpBCr8=
27 github.com/kisielk/gotool v1.0.0/go.mod h1:XhKaO+MFFWcvkIS/tQcRk01m1F5IRFswLeQ+oQHNcck=
28 -github.com/knadh/koanf/maps v0.1.1 h1:G5TjmUh2D7G2YWf5SQQqSiHRJEjaicvU0KpypqB3NIs=
29 -github.com/knadh/koanf/maps v0.1.1/go.mod h1:npD/QZY3V6ghQDdcQzl1W4ICNVTkohC8E73eI2xW4yI=
30 -github.com/knadh/koanf/providers/confmap v0.1.0 h1:gOkxhHkemwG4LezxxN8DMOFopOPghxRVp7JbIvdvqzU=
31 -github.com/knadh/koanf/providers/confmap v0.1.0/go.mod h1:2uLhxQzJnyHKfxG927awZC7+fyHFdQkd697K4MdLnIU=
32 -github.com/knadh/koanf/v2 v2.1.2 h1:I2rtLRqXRy1p01m/utEtpZSSA6dcJbgGVuE27kW2PzQ=
33 -github.com/knadh/koanf/v2 v2.1.2/go.mod h1:Gphfaen0q1Fc1HTgJgSTC4oRX9R2R5ErYMZJy8fLJBo=
28 +github.com/knadh/koanf/maps v0.1.2 h1:RBfmAW5CnZT+PJ1CVc1QSJKf4Xu9kxfQgYVQSu8hpbo=
29 +github.com/knadh/koanf/maps v0.1.2/go.mod h1:npD/QZY3V6ghQDdcQzl1W4ICNVTkohC8E73eI2xW4yI=
30 +github.com/knadh/koanf/providers/confmap v1.0.0 h1:mHKLJTE7iXEys6deO5p6olAiZdG5zwp8Aebir+/EaRE=
31 +github.com/knadh/koanf/providers/confmap v1.0.0/go.mod h1:txHYHiI2hAtF0/0sCmcuol4IDcuQbKTybiB1nOcUo1A=
32 +github.com/knadh/koanf/v2 v2.2.0 h1:FZFwd9bUjpb8DyCWARUBy5ovuhDs1lI87dOEn2K8UVU=
33 +github.com/knadh/koanf/v2 v2.2.0/go.mod h1:PSFru3ufQgTsI7IF+95rf9s8XA1+aHxKuO/W+dPoHEY=
34 github.com/kr/pretty v0.3.1 h1:flRD4NNwYAUpkphVc1HcthR4KEIFJ65n8Mw5qdRn3LE=
35 github.com/kr/pretty v0.3.1/go.mod h1:hoEshYVHaxMs3cyo3Yncou5ZscifuDolrwPKZanG3xk=
36 github.com/kr/text v0.2.0 h1:5Nx0Ya0ZqY2ygV366QzturHI13Jq95ApcVaJBhpS+AY=
@@ -56,52 +56,52 @@ github.com/yuin/goldmark v1.1.27/go.mod h1:3hX8gzYuyVAZsxl0MRgGTJEmQBFcNTphYh9de
56 github.com/yuin/goldmark v1.2.1/go.mod h1:3hX8gzYuyVAZsxl0MRgGTJEmQBFcNTphYh9decYSb74=
57 go.opentelemetry.io/auto/sdk v1.1.0 h1:cH53jehLUN6UFLY71z+NDOiNJqDdPRaXzTel0sJySYA=
58 go.opentelemetry.io/auto/sdk v1.1.0/go.mod h1:3wSPjt5PWp2RhlCcmmOial7AvC4DQqZb7a7wCow3W8A=
59 -go.opentelemetry.io/collector/component v1.29.0 h1:SQDKqMb3EUNPyxZDrImkqgi8p682Hm986CMXBiiTkz0=
60 -go.opentelemetry.io/collector/component v1.29.0/go.mod h1:qHoiz0iyBgyAYsT0F8FgtcJv8pTEQDvo1Ox3eO3qVlY=
61 -go.opentelemetry.io/collector/component/componenttest v0.123.0 h1:h0B/kBj0URKq+i9iMbMuLhc6/dZ2GWL0y9L6tqHRNuA=
62 -go.opentelemetry.io/collector/component/componenttest v0.123.0/go.mod h1:4Y6EMvsgE9fUNM98G0eW5+LFXfcxdhTHQDhaOxJrgN8=
63 -go.opentelemetry.io/collector/config/configretry v1.29.0 h1:Gm140dYwQk9TRQBQHF4hYybdvE2KvUhCLndPvZupTxM=
64 -go.opentelemetry.io/collector/config/configretry v1.29.0/go.mod h1:QNnb+MCk7aS1k2EuGJMtlNCltzD7b8uC7Xel0Dxm1wQ=
65 -go.opentelemetry.io/collector/confmap v1.29.0 h1:1d43r3gQApgRxjyiQ86/qQd7OHW/BEbM6m/L36O5gKU=
66 -go.opentelemetry.io/collector/confmap v1.29.0/go.mod h1:g8DLibnb2MEUWMMlyxHl1VviwFKUdwyBE5M1bxY1Pwg=
67 -go.opentelemetry.io/collector/consumer v1.29.0 h1:WOtwc7bY0/L47mZPE6fu/txD6Fg2dPI3c/Cq+H1RRNY=
68 -go.opentelemetry.io/collector/consumer v1.29.0/go.mod h1:SRcOjDfuvTEk6pchJ0HuOBLsmbRG26+mZLrwXb9T3TQ=
69 -go.opentelemetry.io/collector/consumer/consumererror v0.123.0 h1:/X32KeSGTLX+/YRlR1Wx6/GQkL2m1A6l3B0BYErnxC0=
70 -go.opentelemetry.io/collector/consumer/consumererror v0.123.0/go.mod h1:A9oAkn1emXnenyXiBb5iJNteWCidrM5remnzRqXHOKc=
71 -go.opentelemetry.io/collector/consumer/consumertest v0.123.0 h1:IGBan9tb/QJOzGO8AaQai/jv2DQ7V1DbsFhrTHhDSe8=
72 -go.opentelemetry.io/collector/consumer/consumertest v0.123.0/go.mod h1:MN6OZWsjF1JpvxEzGy7xFu0+PGtbzOASQnN6GgX7olA=
73 -go.opentelemetry.io/collector/consumer/xconsumer v0.123.0 h1:UJhwYgtAe3O/8ilQj+AiyLBz+PXnw/hrpPk3wkB7u9g=
74 -go.opentelemetry.io/collector/consumer/xconsumer v0.123.0/go.mod h1:+LAk3gs6kQ3O0+jLZhZAfgPEWUinxKkFqb3nZkYbC6w=
75 -go.opentelemetry.io/collector/exporter v0.123.0 h1:Dda4AJl6GB+3s6syq45EGpP/kq3xnjcYBBd9SL5A0iA=
76 -go.opentelemetry.io/collector/exporter v0.123.0/go.mod h1:ZeRhRP7AfWmkWkQIBn6A9sIzPkhpjFSjC6Cr78kugr0=
77 -go.opentelemetry.io/collector/exporter/exportertest v0.123.0 h1:TR8pmnyUvy6BnF/suwplZ8l/s7ozetlE1sRYeNfvpoc=
78 -go.opentelemetry.io/collector/exporter/exportertest v0.123.0/go.mod h1:GHt8/K/1evqJ84IUT1QjYrpmBS5dtQOjc7moxik7rA0=
79 -go.opentelemetry.io/collector/exporter/xexporter v0.123.0 h1:7uAARN1ktH8OnhansCiJkhqYcfjp4UtdjUimeORJCYc=
80 -go.opentelemetry.io/collector/exporter/xexporter v0.123.0/go.mod h1:k8TLjOKKm7RxoXIk/aA7CF7CdsoCFYx+4Sfnjd2rR5U=
81 -go.opentelemetry.io/collector/extension v1.29.0 h1:xydxgo+UwDSZtcYGjEbj2UJtfevBSqRhPYyDuOC2VSI=
82 -go.opentelemetry.io/collector/extension v1.29.0/go.mod h1:s/8CiRG1UdAEmovkO7zGACoxBqhqo0UaHciK4Ud+5P4=
83 -go.opentelemetry.io/collector/extension/extensiontest v0.123.0 h1:QankUyQPcmx0UIBGWJX8ObpdESHULhCo0JDhXtagb0s=
84 -go.opentelemetry.io/collector/extension/extensiontest v0.123.0/go.mod h1:JSnbeMyH55jFQqT77wZh04Xx8EnwcIY39A+MqnwDawk=
85 -go.opentelemetry.io/collector/extension/xextension v0.123.0 h1:GA0Hp9+3aY+bCXlPTsSSkQKgRpK3TMJ7lke19DGyqXA=
86 -go.opentelemetry.io/collector/extension/xextension v0.123.0/go.mod h1:qMBSQaxnFDL74poKWeVifhhuF5UCZ73GGD0mgyTI6SE=
87 -go.opentelemetry.io/collector/featuregate v1.29.0 h1:DTDeaokT53AA+SiYBH8GC893cmppThU65b3XyFfLPWw=
88 -go.opentelemetry.io/collector/featuregate v1.29.0/go.mod h1:Y/KsHbvREENKvvN9RlpiWk/IGBK+CATBYzIIpU7nccc=
89 -go.opentelemetry.io/collector/internal/telemetry v0.123.0 h1:K6UMv22ulxcv/G8BkwywMN4LprpwuLl1rhTMWpnujqE=
90 -go.opentelemetry.io/collector/internal/telemetry v0.123.0/go.mod h1:Qwh5XcqAlAT9NorjaaSaNcilDgaecs8S+OKljJ9hFno=
91 -go.opentelemetry.io/collector/pdata v1.29.0 h1:ZXSZ2fROdAEbv4JKKiCspBpjIjYZ5XaNt71LNH4RpQw=
92 -go.opentelemetry.io/collector/pdata v1.29.0/go.mod h1:9kb3zMtLFXBPA6WGWkBHbkFwlwwYL/OHk1m0ASWZpeY=
93 -go.opentelemetry.io/collector/pdata/pprofile v0.123.0 h1:NjwE1btsQbzOWsIFYt7vfgrqDNCROO3hWiZwT1/KaZU=
94 -go.opentelemetry.io/collector/pdata/pprofile v0.123.0/go.mod h1:af2WDdbCc+JA+FokOnsW3l+agsnW5ed8LkLn6CfneII=
95 -go.opentelemetry.io/collector/pdata/testdata v0.123.0 h1:daMcuywv7L/O0pJMCKzlgsm/En8o7SeMSPWUZsPt+JM=
96 -go.opentelemetry.io/collector/pdata/testdata v0.123.0/go.mod h1:lyARYT+2evYj7EdB0LLuo81lb48P7/yKKT86lRCEsmA=
97 -go.opentelemetry.io/collector/pipeline v0.123.0 h1:LDcuCrwhCTx2yROJZqhNmq2v0CFkCkUEvxvvcRW0+2c=
98 -go.opentelemetry.io/collector/pipeline v0.123.0/go.mod h1:TO02zju/K6E+oFIOdi372Wk0MXd+Szy72zcTsFQwXl4=
99 -go.opentelemetry.io/collector/receiver v1.29.0 h1:HK0bzpRlcqqdWBG2v/N76f0eauN8woIU91LloLLWeGM=
100 -go.opentelemetry.io/collector/receiver v1.29.0/go.mod h1:4QoZQZl4r9bIrKa7VLRw3Rp7++ezc5QoXn9VoPZwfoI=
101 -go.opentelemetry.io/collector/receiver/receivertest v0.123.0 h1:MB2x4v7Pg+npLyIQxBn/7rrW4RliV1JtcKqWM75iz1U=
102 -go.opentelemetry.io/collector/receiver/receivertest v0.123.0/go.mod h1:eDcm8N4yFz//V5dSa6MeEdjuK5Nq1X+j9xP9kPJr8BU=
103 -go.opentelemetry.io/collector/receiver/xreceiver v0.123.0 h1:oRExkR/d5axuvAUa3lcOf1ea+zL0hK+iGAhewqZz1NU=
104 -go.opentelemetry.io/collector/receiver/xreceiver v0.123.0/go.mod h1:M3a0ACBOPHag90EV+GbIWKZhJkHXW7Ka2xabBGY25aw=
59 +go.opentelemetry.io/collector/component v1.30.0 h1:HXjqBHaQ47/EEuWdnkjr4Y3kRWvmyWIDvqa1Q262Fls=
60 +go.opentelemetry.io/collector/component v1.30.0/go.mod h1:vfM9kN+BM6oHBXWibquiprz8CVawxd4/aYy3nbhme3E=
61 +go.opentelemetry.io/collector/component/componenttest v0.124.0 h1:Wsc+DmDrWTFs/aEyjDA3slNwV+h/0NOyIR5Aywvr6Zw=
62 +go.opentelemetry.io/collector/component/componenttest v0.124.0/go.mod h1:NQ4ATOzMFc7QA06B993tq8o27DR0cu/JR/zK7slGJ3E=
63 +go.opentelemetry.io/collector/config/configretry v1.30.0 h1:sapni1tymwNiuI0PjqlRR5CvYxIQYT8tyjQGVJDkVPM=
64 +go.opentelemetry.io/collector/config/configretry v1.30.0/go.mod h1:QNnb+MCk7aS1k2EuGJMtlNCltzD7b8uC7Xel0Dxm1wQ=
65 +go.opentelemetry.io/collector/confmap v1.30.0 h1:Y0MXhjQCdMyJN9xZMWWdNPWs6ncMVf7YVnyAEN2dAcM=
66 +go.opentelemetry.io/collector/confmap v1.30.0/go.mod h1:9DdThVDIC3VsdtTb7DgT+HwusWOocoqDkd/TErEtQgA=
67 +go.opentelemetry.io/collector/consumer v1.30.0 h1:Nn6kFTH+EJbv13E0W+sNvWrTgbiFCRv8f6DaA2F1DQs=
68 +go.opentelemetry.io/collector/consumer v1.30.0/go.mod h1:edRyfk61ugdhCQ93PBLRZfYMVWjdMPpKP8z5QLyESf0=
69 +go.opentelemetry.io/collector/consumer/consumererror v0.124.0 h1:OmeJex0C8jcwyILG+eJIGDe6rGaR15fip+Rj3XyMTRY=
70 +go.opentelemetry.io/collector/consumer/consumererror v0.124.0/go.mod h1:A6WL/E70dUh9oFdQsL+aVYGL+d6UqiAhi7rIW/6gvDM=
71 +go.opentelemetry.io/collector/consumer/consumertest v0.124.0 h1:2arChG4RPrHW3lfVWlK/KDF7Y7qkUm/YAiBXh8oTue0=
72 +go.opentelemetry.io/collector/consumer/consumertest v0.124.0/go.mod h1:Hlu+EXbINHxVAyIT1baKO2d0j5odR3fLlLAiaP+JqQg=
73 +go.opentelemetry.io/collector/consumer/xconsumer v0.124.0 h1:/cut96EWVNoz6lIeGI9+EzS6UClMtnZkx5YIpkD0Xe0=
74 +go.opentelemetry.io/collector/consumer/xconsumer v0.124.0/go.mod h1:fHH/MpzFCRNk/4foiYE6BoXQCAMf5sJTO35uvzVrrd4=
75 +go.opentelemetry.io/collector/exporter v0.124.0 h1:ii+9tU/iSrPl4+YDvqFVflksA9hUYEzwMIpmvP4JZ8w=
76 +go.opentelemetry.io/collector/exporter v0.124.0/go.mod h1:Q8tOEwFu3CN8VGjE4H2yZcCRG9Q60foQIyZGKPD/jig=
77 +go.opentelemetry.io/collector/exporter/exportertest v0.124.0 h1:IOxA/4CiVWGPlmA0JofK6W4DzvwW1YJes09r6osluIE=
78 +go.opentelemetry.io/collector/exporter/exportertest v0.124.0/go.mod h1:2EmU8IwVJV79MmFBFFW1LCN0Ob2UZsEkX/mSUB06lbI=
79 +go.opentelemetry.io/collector/exporter/xexporter v0.124.0 h1:Itfn2+F4ki8hObOtPCecWBwGpuxakUYSsTwwkB5iUns=
80 +go.opentelemetry.io/collector/exporter/xexporter v0.124.0/go.mod h1:dNK/PPY02gA9BawIKHyVk8kIFdYvqVZ2A+LlMZucIPY=
81 +go.opentelemetry.io/collector/extension v1.30.0 h1:AJqntAp1p40Q1az2Vze3OHiMURq56KWnUxaLzs1ghaA=
82 +go.opentelemetry.io/collector/extension v1.30.0/go.mod h1:a21WpypFQp9x0Go7yMOknYmIKvdIoWGzjz+h1WMjzLk=
83 +go.opentelemetry.io/collector/extension/extensiontest v0.124.0 h1:pWfKxEqvq5vVdQy+UIerw5j7ZrXezxsfQfuUe8g/xVo=
84 +go.opentelemetry.io/collector/extension/extensiontest v0.124.0/go.mod h1:DLVRyW7tJt8TtYq0Wr5BUsM494YqDiIjN8YCmbVKqjs=
85 +go.opentelemetry.io/collector/extension/xextension v0.124.0 h1:Yzf11HXaiMHfS50Zy/CYKfJjoi+/w/tgRZdDQ2VIdW0=
86 +go.opentelemetry.io/collector/extension/xextension v0.124.0/go.mod h1:GeM0aSgwVSba3Bvvspuy1E+1aa/Q1CDxoK+e/xcJFVg=
87 +go.opentelemetry.io/collector/featuregate v1.30.0 h1:mx7+iP/FQnY7KO8qw/xE3Qd1MQkWcU8VgcqLNrJ8EU8=
88 +go.opentelemetry.io/collector/featuregate v1.30.0/go.mod h1:Y/KsHbvREENKvvN9RlpiWk/IGBK+CATBYzIIpU7nccc=
89 +go.opentelemetry.io/collector/internal/telemetry v0.124.0 h1:kzd1/ZYhLj4bt2pDB529mL4rIRrRacemXodFNxfhdWk=
90 +go.opentelemetry.io/collector/internal/telemetry v0.124.0/go.mod h1:ZjXjqV0dJ+6D4XGhTOxg/WHjnhdmXsmwmUSgALea66Y=
91 +go.opentelemetry.io/collector/pdata v1.30.0 h1:j3jyq9um436r6WzWySzexP2nLnFdmL5uVBYAlyr9nDM=
92 +go.opentelemetry.io/collector/pdata v1.30.0/go.mod h1:0Bxu1ktuj4wE7PIASNSvd0SdBscQ1PLtYasymJ13/Cs=
93 +go.opentelemetry.io/collector/pdata/pprofile v0.124.0 h1:ZjL9wKqzP4BHj0/F1jfGxs1Va8B7xmYayipZeNVoWJE=
94 +go.opentelemetry.io/collector/pdata/pprofile v0.124.0/go.mod h1:1EN3Gw5LSI4fSVma/Yfv/6nqeuYgRTm1/kmG5nE5Oyo=
95 +go.opentelemetry.io/collector/pdata/testdata v0.124.0 h1:vY+pWG7CQfzzGSB5+zGYHQOltRQr59Ek9QiPe+rI+NY=
96 +go.opentelemetry.io/collector/pdata/testdata v0.124.0/go.mod h1:lNH48lGhGv4CYk27fJecpsR1zYHmZjKgNrAprwjym0o=
97 +go.opentelemetry.io/collector/pipeline v0.124.0 h1:hKvhDyH2GPnNO8LGL34ugf36sY7EOXPjBvlrvBhsOdw=
98 +go.opentelemetry.io/collector/pipeline v0.124.0/go.mod h1:TO02zju/K6E+oFIOdi372Wk0MXd+Szy72zcTsFQwXl4=
99 +go.opentelemetry.io/collector/receiver v1.30.0 h1:XbgU4yT3Ld+hL9+jHcD/Kctcr3gXjpiFxKO+50pSayg=
100 +go.opentelemetry.io/collector/receiver v1.30.0/go.mod h1:U3cApz9PHiRMgN0WkZaz4o8mvj1+cVQYsyj2Nl1v3FQ=
101 +go.opentelemetry.io/collector/receiver/receivertest v0.124.0 h1:mx0290aXAo+wfjm4NgbKUodjT5SbS306zmk+AeqeVxE=
102 +go.opentelemetry.io/collector/receiver/receivertest v0.124.0/go.mod h1:3RpopRmIzx5T4zTStHJC0HHfd8YFWm8e9bia1HiuDtY=
103 +go.opentelemetry.io/collector/receiver/xreceiver v0.124.0 h1:YigTUKk8p/aIfqaT0ST7teT9KbLThWD5n2km83byftw=
104 +go.opentelemetry.io/collector/receiver/xreceiver v0.124.0/go.mod h1:NkTpmpAEDT17Dko4gpHUnRztrSkdSd6B0+Y4gfuCWIA=
105 go.opentelemetry.io/contrib/bridges/otelzap v0.10.0 h1:ojdSRDvjrnm30beHOmwsSvLpoRF40MlwNCA+Oo93kXU=
106 go.opentelemetry.io/contrib/bridges/otelzap v0.10.0/go.mod h1:oTTm4g7NEtHSV2i/0FeVdPaPgUIZPfQkFbq0vbzqnv0=
107 go.opentelemetry.io/otel v1.35.0 h1:xKWKPxrxB6OtMCbmMY021CqC45J+3Onta9MqjhnusiQ=
@@ -131,20 +131,20 @@ golang.org/x/net v0.0.0-20190404232315-eb5bcb51f2a3/go.mod h1:t9HGtf8HONx5eT2rtn
131 golang.org/x/net v0.0.0-20190620200207-3b0461eec859/go.mod h1:z5CRVTTTmAJ677TzLLGU+0bjPO0LkuOLi4/5GtJWs/s=
132 golang.org/x/net v0.0.0-20200226121028-0de0cce0169b/go.mod h1:z5CRVTTTmAJ677TzLLGU+0bjPO0LkuOLi4/5GtJWs/s=
133 golang.org/x/net v0.0.0-20201021035429-f5854403a974/go.mod h1:sp8m0HH+o8qH0wwXwYZr8TS3Oi6o0r6Gce1SSxlDquU=
134 -golang.org/x/net v0.37.0 h1:1zLorHbz+LYj7MQlSf1+2tPIIgibq2eL5xkrGk6f+2c=
135 -golang.org/x/net v0.37.0/go.mod h1:ivrbrMbzFq5J41QOQh0siUuly180yBYtLp+CKbEaFx8=
134 +golang.org/x/net v0.39.0 h1:ZCu7HMWDxpXpaiKdhzIfaltL9Lp31x/3fCP11bc6/fY=
135 +golang.org/x/net v0.39.0/go.mod h1:X7NRbYVEA+ewNkCNyJ513WmMdQ3BineSwVtN2zD/d+E=
136 golang.org/x/sync v0.0.0-20190423024810-112230192c58/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM=
137 golang.org/x/sync v0.0.0-20190911185100-cd5d95a43a6e/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM=
138 golang.org/x/sync v0.0.0-20201020160332-67f06af15bc9/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM=
139 golang.org/x/sys v0.0.0-20190215142949-d0b11bdaac8a/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY=
140 golang.org/x/sys v0.0.0-20190412213103-97732733099d/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
141 golang.org/x/sys v0.0.0-20200930185726-fdedc70b468f/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
142 -golang.org/x/sys v0.31.0 h1:ioabZlmFYtWhL+TRYpcnNlLwhyxaM9kWTDEmfnprqik=
143 -golang.org/x/sys v0.31.0/go.mod h1:BJP2sWEmIv4KK5OTEluFJCKSidICx8ciO85XgH3Ak8k=
142 +golang.org/x/sys v0.32.0 h1:s77OFDvIQeibCmezSnk/q6iAfkdiQaJi4VzroCFrN20=
143 +golang.org/x/sys v0.32.0/go.mod h1:BJP2sWEmIv4KK5OTEluFJCKSidICx8ciO85XgH3Ak8k=
144 golang.org/x/text v0.3.0/go.mod h1:NqM8EUOU14njkJ3fqMW+pc6Ldnwhi/IjpwHt7yyuwOQ=
145 golang.org/x/text v0.3.3/go.mod h1:5Zoc/QRtKVWzQhOtBMvqHzDpF6irO9z98xDceosuGiQ=
146 -golang.org/x/text v0.23.0 h1:D71I7dUrlY+VX0gQShAThNGHFxZ13dGLBHQLVl1mJlY=
147 -golang.org/x/text v0.23.0/go.mod h1:/BLNzu4aZCJ1+kcD0DNRotWKage4q2rGVAg4o22unh4=
146 +golang.org/x/text v0.24.0 h1:dd5Bzh4yt5KYA8f9CJHCP4FB4D51c2c6JvN37xJJkJ0=
147 +golang.org/x/text v0.24.0/go.mod h1:L8rBsPeo2pSS+xqN0d5u2ikmjtmoJbDBT1b7nHvFCdU=
148 golang.org/x/tools v0.0.0-20180917221912-90fa682c2a6e/go.mod h1:n7NCudcB/nEzxVGmLbDWY5pfWTLqBcC2KZ6jyYvM4mQ=
149 golang.org/x/tools v0.0.0-20191119224855-298f0cb1881e/go.mod h1:b+2E5dAYhXwXZwtnZ6UAqBI28+e2cm9otk0dWdXHAEo=
150 golang.org/x/tools v0.0.0-20200619180055-7c47624df98f/go.mod h1:EkVYQZoAsY45+roYkvgYkIh4xh/qjgUK9TdY2XT94GE=
@@ -153,10 +153,10 @@ golang.org/x/xerrors v0.0.0-20190717185122-a985d3407aa7/go.mod h1:I/5z698sn9Ka8T
153 golang.org/x/xerrors v0.0.0-20191011141410-1b5146add898/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0=
154 golang.org/x/xerrors v0.0.0-20191204190536-9bdfabe68543/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0=
155 golang.org/x/xerrors v0.0.0-20200804184101-5ec99f83aff1/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0=
156 -google.golang.org/genproto/googleapis/rpc v0.0.0-20250303144028-a0af3efb3deb h1:TLPQVbx1GJ8VKZxz52VAxl1EBgKXXbTiU9Fc5fZeLn4=
157 -google.golang.org/genproto/googleapis/rpc v0.0.0-20250303144028-a0af3efb3deb/go.mod h1:LuRYeWDFV6WOn90g357N17oMCaxpgCnbi/44qJvDn2I=
158 -google.golang.org/grpc v1.71.0 h1:kF77BGdPTQ4/JZWMlb9VpJ5pa25aqvVqogsxNHHdeBg=
159 -google.golang.org/grpc v1.71.0/go.mod h1:H0GRtasmQOh9LkFoCPDu3ZrwUtD1YGE+b2vYBYd/8Ec=
156 +google.golang.org/genproto/googleapis/rpc v0.0.0-20250414145226-207652e42e2e h1:ztQaXfzEXTmCBvbtWYRhJxW+0iJcz2qXfd38/e9l7bA=
157 +google.golang.org/genproto/googleapis/rpc v0.0.0-20250414145226-207652e42e2e/go.mod h1:qQ0YXyHHx3XkvlzUtpXDkS29lDSafHMZBAZDc03LQ3A=
158 +google.golang.org/grpc v1.71.1 h1:ffsFWr7ygTUscGPI0KKK6TLrGz0476KUvvsbqWK0rPI=
159 +google.golang.org/grpc v1.71.1/go.mod h1:H0GRtasmQOh9LkFoCPDu3ZrwUtD1YGE+b2vYBYd/8Ec=
160 google.golang.org/protobuf v1.36.6 h1:z1NpPI8ku2WgiWnf+t9wTPsn6eP1L7ksHUlkfLvd9xY=
161 google.golang.org/protobuf v1.36.6/go.mod h1:jduwjTPXsFjZGTmRluh+L6NjiWu7pchiJ2/5YcXBHnY=
162 gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0=
src/go/otel-collector/exporter/journaldexporter/journal_remote_test.go
+318 -371
@@ -1,471 +1,417 @@
1 package journaldexporter
2
3 import (
4 + "bytes"
5 "context"
6 "errors"
7 "fmt"
8 "io"
9 "net/http"
10 "net/http/httptest"
11 + "strings"
12 "sync"
13 "testing"
14 "time"
15
16 "github.com/stretchr/testify/assert"
17 "github.com/stretchr/testify/require"
18 + "go.uber.org/zap"
19 "go.uber.org/zap/zaptest"
20 )
21
22 func TestRemoteJournalClient_sendMessage(t *testing.T) {
23 tests := map[string]struct {
21 - setupServer func(*testing.T) (*httptest.Server, *chunkWriterHandler)
22 - setupClient func(*testing.T, string) *remoteJournalClient
23 - setupContext func() (context.Context, context.CancelFunc)
24 - messages [][]byte
25 - waitBetween time.Duration
26 - waitAfter time.Duration
27 - validateErr func(*testing.T, error)
28 - validateData func(*testing.T, *chunkWriterHandler, [][]byte)
29 - shutdownFirst bool
24 + prepare func(t *testing.T) (*httptest.Server, *chunkReaderTestHandler)
25 + messages []string
26 + wantErr bool
27 }{
28 "successful message delivery": {
32 - setupServer: prepareTestServer,
33 - setupClient: prepareTestRemoteJournalClient,
34 - setupContext: func() (context.Context, context.CancelFunc) {
35 - return context.Background(), func() {}
29 + prepare: func(t *testing.T) (*httptest.Server, *chunkReaderTestHandler) {
30 + return prepareTestServer(t)
31 },
37 - messages: [][]byte{[]byte("test message data")},
38 - waitAfter: 100 * time.Millisecond,
39 - validateErr: func(t *testing.T, err error) {
40 - assert.NoError(t, err)
41 - },
42 - validateData: func(t *testing.T, handler *chunkWriterHandler, messages [][]byte) {
43 - assert.Equal(t, messages[0], handler.getReceivedData())
32 + messages: []string{
33 + `{"message":"hello world1"}`,
34 + `{"message":"hello world2"}`,
35 + `{"message":"hello world3"}`,
36 },
37 },
46 - "send empty message": {
47 - setupServer: prepareTestServer,
48 - setupClient: prepareTestRemoteJournalClient,
49 - setupContext: func() (context.Context, context.CancelFunc) {
50 - return context.Background(), func() {}
51 - },
52 - messages: [][]byte{{}},
53 - waitAfter: 100 * time.Millisecond,
54 - validateErr: func(t *testing.T, err error) {
55 - assert.NoError(t, err)
56 - },
57 - validateData: func(t *testing.T, handler *chunkWriterHandler, messages [][]byte) {
58 - assert.Empty(t, handler.getReceivedData())
38 + "error message delivery": {
39 + wantErr: true,
40 + prepare: func(t *testing.T) (*httptest.Server, *chunkReaderTestHandler) {
41 + srv, h := prepareTestServer(t)
42 + h.errorAfterRead = true
43 + return srv, h
44 + },
45 + messages: []string{
46 + `{"message":"hello world1"}`,
47 + `{"message":"hello world2"}`,
48 + `{"message":"hello world3"}`,
49 },
50 },
61 - "send message with cancelled context": {
62 - setupServer: prepareTestServer,
63 - setupClient: prepareTestRemoteJournalClient,
64 - setupContext: func() (context.Context, context.CancelFunc) {
65 - ctx, cancel := context.WithCancel(context.Background())
66 - cancel() // Cancel immediately
67 - return ctx, cancel
51 + "connection abrupt close": {
52 + wantErr: true,
53 + prepare: func(t *testing.T) (*httptest.Server, *chunkReaderTestHandler) {
54 + srv, h := prepareTestServer(t)
55 + h.closeEarly = true
56 + return srv, h
57 },
69 - messages: [][]byte{[]byte("test message that should not be sent")},
70 - waitAfter: 100 * time.Millisecond,
71 - validateErr: func(t *testing.T, err error) {
72 - assert.Error(t, err)
73 - assert.True(t, errors.Is(err, context.Canceled))
58 + messages: []string{
59 + `{"message":"hello world1"}`,
60 + `{"message":"hello world2"}`,
61 },
62 },
63 "attempt to send to non-existent server": {
77 - setupServer: nil, // No server
78 - setupClient: func(t *testing.T, _ string) *remoteJournalClient {
79 - return prepareTestRemoteJournalClient(t, "http://non-existent-server:12345")
80 - },
81 - setupContext: func() (context.Context, context.CancelFunc) {
82 - return context.Background(), func() {}
64 + wantErr: true,
65 + prepare: func(t *testing.T) (*httptest.Server, *chunkReaderTestHandler) {
66 + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {}))
67 + srv.Close()
68 +
69 + return srv, nil
70 },
84 - messages: [][]byte{[]byte("test message to non-existent server")},
85 - waitBetween: 100 * time.Millisecond,
86 - validateErr: func(t *testing.T, err error) {
87 - assert.Error(t, err)
88 - //assert.Contains(t, err.Error(), "connection unavailable")
71 + messages: []string{
72 + `{"message":"hello world1"}`,
73 },
74 },
91 - "send after client is shut down": {
92 - setupServer: prepareTestServer,
93 - setupClient: prepareTestRemoteJournalClient,
94 - setupContext: func() (context.Context, context.CancelFunc) {
95 - return context.Background(), func() {}
75 + "empty message": {
76 + prepare: func(t *testing.T) (*httptest.Server, *chunkReaderTestHandler) {
77 + return prepareTestServer(t)
78 },
97 - messages: [][]byte{[]byte("test message after shutdown")},
98 - shutdownFirst: true,
99 - validateErr: func(t *testing.T, err error) {
100 - assert.Error(t, err)
101 - assert.Contains(t, err.Error(), "client is shut down")
79 + messages: []string{
80 + ``,
81 },
82 },
104 - // TODO: fix
105 - //"server responds with error status": {
106 - // setupServer: func(t *testing.T) (*httptest.Server, *chunkWriterHandler) {
107 - // handler := &chunkWriterHandler{
108 - // statusCode: http.StatusInternalServerError,
109 - // }
110 - // server := httptest.NewServer(handler)
111 - // t.Cleanup(func() {
112 - // server.Close()
113 - // })
114 - // return server, handler
115 - // },
116 - // setupClient: prepareTestRemoteJournalClient,
117 - // setupContext: func() (context.Context, context.CancelFunc) {
118 - // return context.Background(), func() {}
119 - // },
120 - // messages: [][]byte{
121 - // []byte("test message expecting error"),
122 - // []byte("second message after error"),
123 - // },
124 - // waitBetween: 300 * time.Millisecond, // Longer wait to ensure background error is processed
125 - // waitAfter: 100 * time.Millisecond,
126 - // validateErr: func(t *testing.T, err error) {
127 - // // The first message might succeed because the error only occurs after the server responds
128 - // // But the second message should fail once the background goroutine sets reqErr
129 - // require.Error(t, err)
130 - // // If we get an error about pipe being closed or connection unavailable, that's valid
131 - // valid := strings.Contains(err.Error(), "connection unavailable") ||
132 - // strings.Contains(err.Error(), "failed to write message") ||
133 - // strings.Contains(err.Error(), "pipe is closed")
134 - // assert.True(t, valid, "Expected error to indicate connection problem, got: %v", err)
135 - // },
136 - //},
137 - // TODO: fix
138 - //"server closes connection mid-stream": {
139 - // setupServer: func(t *testing.T) (*httptest.Server, *chunkWriterHandler) {
140 - // handler := &chunkWriterHandler{
141 - // statusCode: http.StatusOK,
142 - // closeEarly: true,
143 - // }
144 - // server := httptest.NewServer(handler)
145 - // t.Cleanup(func() {
146 - // server.Close()
147 - // })
148 - // return server, handler
149 - // },
150 - // setupClient: prepareTestRemoteJournalClient,
151 - // setupContext: func() (context.Context, context.CancelFunc) {
152 - // return context.Background(), func() {}
153 - // },
154 - // messages: [][]byte{[]byte("first message that will cause server to close"), []byte("second message after closure")},
155 - // waitBetween: 200 * time.Millisecond,
156 - // validateErr: func(t *testing.T, err error) {
157 - // assert.Error(t, err)
158 - // assert.Contains(t, err.Error(), "connection unavailable")
159 - // },
160 - //},
161 - "sending multiple messages": {
162 - setupServer: prepareTestServer,
163 - setupClient: prepareTestRemoteJournalClient,
164 - setupContext: func() (context.Context, context.CancelFunc) {
165 - return context.Background(), func() {}
83 + "large message": {
84 + prepare: func(t *testing.T) (*httptest.Server, *chunkReaderTestHandler) {
85 + return prepareTestServer(t)
86 },
167 - messages: [][]byte{
168 - []byte("first message"),
169 - []byte("second message"),
170 - []byte("third message"),
87 + messages: []string{
88 + strings.Repeat(`{"message":"large payload test"}`, 1000),
89 },
172 - waitAfter: 100 * time.Millisecond,
173 - validateErr: func(t *testing.T, err error) {
174 - assert.NoError(t, err)
90 + },
91 + "slow server response": {
92 + prepare: func(t *testing.T) (*httptest.Server, *chunkReaderTestHandler) {
93 + srv, h := prepareTestServer(t)
94 + h.readDelay = 200 * time.Millisecond
95 + return srv, h
96 },
176 - validateData: func(t *testing.T, handler *chunkWriterHandler, messages [][]byte) {
177 - expected := append(append(messages[0], messages[1]...), messages[2]...)
178 - assert.Equal(t, expected, handler.getReceivedData())
97 + messages: []string{
98 + `{"message":"hello world1"}`,
99 + `{"message":"hello world2"}`,
100 },
101 },
102 }
103
183 - for name, tc := range tests {
104 + for name, test := range tests {
105 t.Run(name, func(t *testing.T) {
185 - var server *httptest.Server
186 - var handler *chunkWriterHandler
187 - var client *remoteJournalClient
106 + srv, h := test.prepare(t)
107 + if srv != nil {
108 + defer srv.Close()
109 + }
110
189 - if tc.setupServer != nil {
190 - server, handler = tc.setupServer(t)
191 - client = tc.setupClient(t, server.URL)
111 + jc := prepareRemoteJournalClient(t, srv.URL)
112 +
113 + var wg sync.WaitGroup
114 + sendErrCh := make(chan error, 1)
115 + recvErrCh := make(chan error, 1)
116 +
117 + // Skip receive handling for non-existent server and empty message cases
118 + skipReceive := h == nil || len(test.messages) == 0 ||
119 + (len(test.messages) == 1 && test.messages[0] == "")
120 +
121 + wg.Add(1)
122 + go func() {
123 + defer wg.Done()
124 + var err error
125 + for _, msg := range test.messages {
126 + err = errors.Join(err, jc.sendMessage(context.Background(), []byte(msg)))
127 + // Brief pause to allow processing
128 + time.Sleep(50 * time.Millisecond)
129 + }
130 + sendErrCh <- err
131 + }()
132 +
133 + if skipReceive {
134 + // For non-server or empty message tests, we don't expect to receive anything
135 + recvErrCh <- nil
136 } else {
193 - client = tc.setupClient(t, "")
137 + wg.Add(1)
138 + go func() {
139 + defer wg.Done()
140 +
141 + want := strings.Join(test.messages, "")
142 + var buf bytes.Buffer
143 + timeout := 5 * time.Second
144 +
145 + // For large payloads, extend timeout
146 + if len(want) > 10000 {
147 + timeout = 10 * time.Second
148 + }
149 +
150 + timeoutCh := time.After(timeout)
151 +
152 + for {
153 + select {
154 + case msg, ok := <-h.recDataCh:
155 + if !ok {
156 + // If we're expecting an error due to connection close,
157 + // this is normal and not a test failure
158 + if test.wantErr && buf.Len() > 0 {
159 + recvErrCh <- nil
160 + } else if buf.String() == want {
161 + recvErrCh <- nil
162 + } else if buf.Len() > 0 {
163 + recvErrCh <- fmt.Errorf("received incomplete data: got %q, want %q",
164 + buf.String(), want)
165 + } else {
166 + recvErrCh <- errors.New("channel closed without receiving data")
167 + }
168 + return
169 + }
170 + _, _ = buf.Write(bytes.Clone(msg))
171 + if buf.String() == want {
172 + recvErrCh <- nil
173 + return
174 + }
175 + case <-timeoutCh:
176 + recvErrCh <- fmt.Errorf("timeout waiting for messages after %v", timeout)
177 + return
178 + }
179 + }
180 + }()
181 }
195 - defer func() { _ = client.shutdown(context.Background()) }()
182
197 - ctx, cancel := tc.setupContext()
198 - defer cancel()
183 + done := make(chan struct{})
184 + go func() {
185 + defer close(done)
186 + wg.Wait()
187 + }()
188 +
189 + select {
190 + case <-done:
191 + case <-time.After(10 * time.Second):
192 + _ = jc.shutdown(context.Background())
193 + t.Fatal("timed out waiting for test completion")
194 + }
195
200 - if tc.shutdownFirst {
201 - err := client.shutdown(context.Background())
202 - assert.NoError(t, err)
196 + shutdownCtx, cancel := context.WithTimeout(context.Background(), 2*time.Second)
197 + defer cancel()
198 + err := jc.shutdown(shutdownCtx)
199 + if err != nil {
200 + t.Logf("Shutdown error (non-fatal): %v", err)
201 }
202
205 - var lastErr error
206 - for i, msg := range tc.messages {
207 - lastErr = client.sendMessage(ctx, msg)
203 + // Give the test a moment to complete any pending operations
204 + time.Sleep(100 * time.Millisecond)
205
209 - // For multi-message tests, we only validate error on the last message
210 - if i < len(tc.messages)-1 && tc.waitBetween > 0 {
211 - time.Sleep(tc.waitBetween)
212 - }
206 + var sendErr, recvErr error
207 + select {
208 + case sendErr = <-sendErrCh:
209 + default:
210 + sendErr = errors.New("send goroutine did not complete")
211 }
212
215 - if tc.waitAfter > 0 {
216 - time.Sleep(tc.waitAfter)
213 + if !skipReceive {
214 + select {
215 + case recvErr = <-recvErrCh:
216 + default:
217 + recvErr = errors.New("receive goroutine did not complete")
218 + }
219 }
220
219 - if tc.validateErr != nil {
220 - tc.validateErr(t, lastErr)
221 - }
221 + combinedErr := errors.Join(sendErr, recvErr)
222
223 - if tc.validateData != nil && handler != nil {
224 - tc.validateData(t, handler, tc.messages)
223 + if test.wantErr {
224 + if name == "empty message" {
225 + // Empty message should not cause an error
226 + assert.NoError(t, combinedErr, "empty message should not cause error")
227 + } else {
228 + assert.Error(t, combinedErr, "expected an error but got none")
229 + }
230 + } else {
231 + assert.NoError(t, combinedErr, "unexpected error")
232 }
233 })
234 }
235 }
236
230 -func TestRemoteJournalClient_shutdown(t *testing.T) {
231 - tests := map[string]struct {
232 - setupServer func(*testing.T) (*httptest.Server, *chunkWriterHandler)
233 - sendMessages func(*testing.T, context.Context, *remoteJournalClient)
234 - setupContext func() (context.Context, context.CancelFunc)
235 - validateErr func(*testing.T, error)
236 - validateShutdown func(*testing.T, *remoteJournalClient, *chunkWriterHandler)
237 - doubleShutdown bool
238 - }{
239 - "normal shutdown": {
240 - setupServer: prepareTestServer,
241 - sendMessages: func(t *testing.T, ctx context.Context, client *remoteJournalClient) {
242 - msg := []byte("test message before shutdown")
243 - err := client.sendMessage(ctx, msg)
244 - assert.NoError(t, err)
245 - },
246 - setupContext: func() (context.Context, context.CancelFunc) {
247 - return context.Background(), func() {}
248 - },
249 - validateErr: func(t *testing.T, err error) {
250 - assert.NoError(t, err)
251 - },
252 - validateShutdown: func(t *testing.T, client *remoteJournalClient, handler *chunkWriterHandler) {
253 - select {
254 - case <-client.done:
255 - // Expected - channel should be closed
256 - default:
257 - t.Error("Expected done channel to be closed after shutdown")
237 +func TestRemoteJournalClient_ConcurrentMessages(t *testing.T) {
238 + srv, h := prepareTestServer(t)
239 + defer srv.Close()
240 +
241 + jc := prepareRemoteJournalClient(t, srv.URL)
242 +
243 + const numMessages = 10
244 + const numGoroutines = 5
245 +
246 + var wg sync.WaitGroup
247 + for i := 0; i < numGoroutines; i++ {
248 + wg.Add(1)
249 + go func(id int) {
250 + defer wg.Done()
251 + for j := 0; j < numMessages; j++ {
252 + msg := fmt.Sprintf(`{"goroutine":%d,"message":%d}`, id, j)
253 + err := jc.sendMessage(context.Background(), []byte(msg))
254 + if err != nil {
255 + t.Errorf("Error sending message from goroutine %d: %v", id, err)
256 + return
257 }
258 + // Small delay to interleave messages
259 + time.Sleep(10 * time.Millisecond)
260 + }
261 + }(i)
262 + }
263
260 - err := client.sendMessage(context.Background(), []byte("test message after shutdown"))
261 - assert.Error(t, err)
262 - assert.Contains(t, err.Error(), "client is shut down")
263 - },
264 - },
265 - // TODO: fix
266 - //"shutdown with timeout": {
267 - // setupServer: func(t *testing.T) (*httptest.Server, *chunkWriterHandler) {
268 - // handler := &chunkWriterHandler{
269 - // statusCode: http.StatusOK,
270 - // delay: 500 * time.Millisecond, // Delay each read to simulate slow processing
271 - // }
272 - // server := httptest.NewServer(handler)
273 - // t.Cleanup(func() {
274 - // server.Close()
275 - // })
276 - // return server, handler
277 - // },
278 - // sendMessages: func(t *testing.T, ctx context.Context, client *remoteJournalClient) {
279 - // // Send a large message to ensure the background task will be busy
280 - // msg := []byte(strings.Repeat("test message with shutdown timeout ", 1000))
281 - // err := client.sendMessage(ctx, msg)
282 - // assert.NoError(t, err)
283 - // },
284 - // setupContext: func() (context.Context, context.CancelFunc) {
285 - // return context.WithTimeout(context.Background(), 100*time.Millisecond)
286 - // },
287 - // validateErr: func(t *testing.T, err error) {
288 - // assert.Error(t, err)
289 - // assert.True(t, errors.Is(err, context.DeadlineExceeded))
290 - // },
291 - // validateShutdown: func(t *testing.T, client *remoteJournalClient, handler *chunkWriterHandler) {
292 - // select {
293 - // case <-client.done:
294 - // // Expected - channel should be closed
295 - // default:
296 - // t.Error("Expected done channel to be closed after shutdown timeout")
297 - // }
298 - // },
299 - //},
300 - "double shutdown": {
301 - setupServer: prepareTestServer,
302 - sendMessages: func(t *testing.T, ctx context.Context, client *remoteJournalClient) {
303 - // No messages sent
304 - },
305 - setupContext: func() (context.Context, context.CancelFunc) {
306 - return context.Background(), func() {}
307 - },
308 - validateErr: func(t *testing.T, err error) {
309 - assert.NoError(t, err)
310 - },
311 - doubleShutdown: true,
312 - },
313 - "shutdown with active connections": {
314 - setupServer: prepareTestServer,
315 - sendMessages: func(t *testing.T, ctx context.Context, client *remoteJournalClient) {
316 - // Send multiple messages to ensure there's an active connection
317 - for i := 0; i < 5; i++ {
318 - msg := []byte(fmt.Sprintf("test message %d before shutdown", i))
319 - err := client.sendMessage(ctx, msg)
320 - assert.NoError(t, err)
321 - }
322 - // Allow some time for the background goroutine to process
323 - time.Sleep(50 * time.Millisecond)
324 - },
325 - setupContext: func() (context.Context, context.CancelFunc) {
326 - return context.Background(), func() {}
327 - },
328 - validateErr: func(t *testing.T, err error) {
329 - assert.NoError(t, err)
330 - },
331 - validateShutdown: func(t *testing.T, client *remoteJournalClient, handler *chunkWriterHandler) {
332 - // Verify at least some data was received before shutdown
333 - receivedData := handler.getReceivedData()
334 - assert.NotEmpty(t, receivedData)
335 - assert.Contains(t, string(receivedData), "test message")
336 - },
337 - },
338 - "shutdown without active connections": {
339 - setupServer: prepareTestServer,
340 - sendMessages: func(t *testing.T, ctx context.Context, client *remoteJournalClient) {
341 - // No messages sent
342 - },
343 - setupContext: func() (context.Context, context.CancelFunc) {
344 - return context.Background(), func() {}
345 - },
346 - validateErr: func(t *testing.T, err error) {
347 - assert.NoError(t, err)
348 - },
349 - validateShutdown: func(t *testing.T, client *remoteJournalClient, handler *chunkWriterHandler) {
350 - select {
351 - case <-client.done:
352 - // Expected - channel should be closed
353 - default:
354 - t.Error("Expected done channel to be closed after shutdown")
264 + received := make([]string, 0, numMessages*numGoroutines)
265 + receiveDone := make(chan struct{})
266 +
267 + go func() {
268 + defer close(receiveDone)
269 + for {
270 + select {
271 + case msg, ok := <-h.recDataCh:
272 + if !ok {
273 + return
274 }
356 - },
357 - },
275 + received = append(received, string(msg))
276 + case <-time.After(5 * time.Second):
277 + t.Error("Timeout waiting for messages")
278 + return
279 + }
280 + }
281 + }()
282 +
283 + wg.Wait()
284 +
285 + err := jc.shutdown(context.Background())
286 + require.NoError(t, err)
287 +
288 + select {
289 + case <-receiveDone:
290 + case <-time.After(2 * time.Second):
291 + t.Fatal("Timeout waiting for receiver to complete")
292 }
293
360 - for name, tc := range tests {
361 - t.Run(name, func(t *testing.T) {
362 - server, handler := tc.setupServer(t)
294 + t.Logf("Received %d message chunks", len(received))
295 + fullMessage := strings.Join(received, "")
296 +
297 + messageCount := strings.Count(fullMessage, `{"goroutine":`)
298 + assert.GreaterOrEqual(t, messageCount, numMessages*numGoroutines,
299 + "Did not receive all expected messages")
300 +}
301
364 - client := prepareTestRemoteJournalClient(t, server.URL)
302 +func TestRemoteJournalClient_ContextCancellation(t *testing.T) {
303 + srv, _ := prepareTestServer(t)
304 + defer srv.Close()
305
366 - // Setup context for messages and initial operations
367 - msgCtx, msgCancel := context.WithTimeout(context.Background(), 5*time.Second)
368 - defer msgCancel()
306 + jc := prepareRemoteJournalClient(t, srv.URL)
307
370 - tc.sendMessages(t, msgCtx, client)
308 + ctx, cancel := context.WithCancel(context.Background())
309
372 - shutdownCtx, shutdownCancel := tc.setupContext()
373 - defer shutdownCancel()
310 + cancel()
311
375 - err := client.shutdown(shutdownCtx)
312 + err := jc.sendMessage(ctx, []byte(`{"message":"should not be sent"}`))
313 + assert.Error(t, err, "Expected error due to cancelled context")
314 + assert.ErrorIs(t, err, context.Canceled, "Error should be context.Canceled")
315
377 - if tc.doubleShutdown {
378 - secondErr := client.shutdown(context.Background())
379 - assert.NoError(t, secondErr, "Second shutdown should not error")
380 - }
316 + _ = jc.shutdown(context.Background())
317 +}
318
382 - if tc.validateErr != nil {
383 - tc.validateErr(t, err)
384 - }
319 +func TestRemoteJournalClient_SendAfterShutdown(t *testing.T) {
320 + srv, _ := prepareTestServer(t)
321 + defer srv.Close()
322
386 - if tc.validateShutdown != nil {
387 - tc.validateShutdown(t, client, handler)
388 - }
389 - })
390 - }
323 + jc := prepareRemoteJournalClient(t, srv.URL)
324 +
325 + err := jc.shutdown(context.Background())
326 + require.NoError(t, err, "Shutdown should succeed")
327 +
328 + err = jc.sendMessage(context.Background(), []byte(`{"message":"after shutdown"}`))
329 + assert.Error(t, err, "Expected error when sending to shut down client")
330 + assert.Contains(t, err.Error(), "client is shut down",
331 + "Error should indicate client is shut down")
332 +}
333 +
334 +type chunkReaderTestHandler struct {
335 + log *zap.Logger
336 + recDataCh chan []byte
337 + readDelay time.Duration // Delay between reads (simulates slow processing)
338 + errorAfterRead bool // If true, send an error status *after* reading at least one chunk
339 + closeEarly bool // If true, simulate an abrupt connection close *after* reading at least one chunk
340 }
341
393 -type chunkWriterHandler struct {
394 - receivedData []byte
395 - mu sync.Mutex
396 - delay time.Duration
397 - statusCode int
398 - errorAfterRead bool
399 - closeEarly bool
342 +func newChunkReaderTestHandler(t *testing.T) *chunkReaderTestHandler {
343 + return &chunkReaderTestHandler{
344 + log: zaptest.NewLogger(t),
345 + recDataCh: make(chan []byte),
346 + }
347 }
348
402 -func (h *chunkWriterHandler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
403 - h.mu.Lock()
404 - h.receivedData = make([]byte, 0)
405 - h.mu.Unlock()
349 +func (h *chunkReaderTestHandler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
350 + defer func() {
351 + close(h.recDataCh)
352 + }()
353
354 contentType := r.Header.Get("Content-Type")
355 if contentType != "application/vnd.fdo.journal" {
409 - w.WriteHeader(http.StatusBadRequest)
410 - _, _ = w.Write([]byte("Invalid content type"))
356 + h.log.Error("Received invalid content type", zap.String("content-type", contentType))
357 + http.Error(w, "Invalid content type", http.StatusBadRequest)
358 return
359 }
360
414 - w.Header().Set("Transfer-Encoding", "chunked")
415 - w.WriteHeader(h.statusCode)
416 -
361 buf := make([]byte, 1024)
418 - readCount := 0
362
420 - readDone := make(chan struct{})
421 - go func() {
422 - defer close(readDone)
363 + for {
364 + if h.readDelay > 0 {
365 + time.Sleep(h.readDelay)
366 + }
367
424 - for {
425 - if h.delay > 0 {
426 - time.Sleep(h.delay)
427 - }
428 - n, err := r.Body.Read(buf)
429 - if n > 0 {
430 - h.mu.Lock()
431 - h.receivedData = append(h.receivedData, buf[:n]...)
432 - h.mu.Unlock()
433 - readCount++
434 - }
368 + n, err := r.Body.Read(buf)
369 + if err != nil {
370 if err == io.EOF {
371 + h.log.Info("Finished reading request body (EOF)")
372 break
373 }
438 - if err != nil {
439 - return
440 - }
441 - if h.errorAfterRead && readCount >= 1 {
442 - if flusher, ok := w.(http.Flusher); ok {
443 - flusher.Flush()
444 - }
445 - return
446 - }
447 - if h.closeEarly && readCount >= 1 {
448 - return
449 - }
450 - if h.statusCode < 200 || h.statusCode >= 300 {
374 + h.log.Error("Error reading request body", zap.Error(err))
375 + http.Error(w, fmt.Sprintf("Error reading request body: %v", err), http.StatusInternalServerError)
376 + return
377 + }
378 +
379 + if h.errorAfterRead {
380 + h.log.Info("Simulating error response after read")
381 + http.Error(w, "Simulated processing error after read", http.StatusInternalServerError)
382 + return
383 + }
384 +
385 + if h.closeEarly {
386 + h.log.Info("Simulating connection abrupt close")
387 + hj, ok := w.(http.Hijacker)
388 + if !ok {
389 + h.log.Warn("Cannot simulate close early: Hijacking not supported")
390 + http.Error(w, "Cannot simulate early close (hijacking not supported)", http.StatusInternalServerError)
391 return
392 }
393 + conn, _, _ := hj.Hijack()
394 + _ = conn.Close()
395 + fmt.Println("Connection closed abruptly.")
396 + return
397 }
454 - }()
398
456 - <-readDone
457 -}
399 + select {
400 + case h.recDataCh <- bytes.Clone(buf[:n]):
401 + case <-time.After(time.Second * 5):
402 + http.Error(w, "Sending read data back timed out", http.StatusLocked)
403 + return
404 + }
405 + }
406
459 -func (h *chunkWriterHandler) getReceivedData() []byte {
460 - h.mu.Lock()
461 - defer h.mu.Unlock()
462 - return h.receivedData
407 + _, _ = io.Copy(io.Discard, r.Body)
408 + _ = r.Body.Close()
409 + w.WriteHeader(http.StatusOK)
410 }
411
465 -func prepareTestServer(t *testing.T) (*httptest.Server, *chunkWriterHandler) {
466 - handler := &chunkWriterHandler{
467 - statusCode: http.StatusOK,
468 - }
412 +func prepareTestServer(t *testing.T) (*httptest.Server, *chunkReaderTestHandler) {
413 + handler := newChunkReaderTestHandler(t)
414 +
415 server := httptest.NewServer(handler)
416 t.Cleanup(func() {
417 server.Close()
@@ -473,10 +419,10 @@ func prepareTestServer(t *testing.T) (*httptest.Server, *chunkWriterHandler) {
419 return server, handler
420 }
421
476 -func prepareTestRemoteJournalClient(t *testing.T, serverURL string) *remoteJournalClient {
422 +func prepareRemoteJournalClient(t *testing.T, serverURL string) *remoteJournalClient {
423 cfg := &Config{
424 URL: serverURL,
479 - Timeout: 1 * time.Second,
425 + Timeout: 3 * time.Second,
426 }
427
428 client, err := newRemoteJournalClient(cfg)
@@ -484,5 +430,6 @@ func prepareTestRemoteJournalClient(t *testing.T, serverURL string) *remoteJourn
430 require.NotNil(t, client)
431
432 client.log = zaptest.NewLogger(t)
433 +
434 return client
435 }
src/go/otel-collector/release-config.yaml.in
+5 -5
@@ -8,13 +8,13 @@ dist:
8 debug_compilation: @DEBUG_BUILD@
9
10 receivers:
11 - - gomod: go.opentelemetry.io/collector/receiver/otlpreceiver v0.123.0
12 - - gomod: github.com/open-telemetry/opentelemetry-collector-contrib/receiver/filelogreceiver v0.123.0
11 + - gomod: go.opentelemetry.io/collector/receiver/otlpreceiver v0.124.0
12 + - gomod: github.com/open-telemetry/opentelemetry-collector-contrib/receiver/filelogreceiver v0.124.0
13
14 exporters:
15 - - gomod: go.opentelemetry.io/collector/exporter/debugexporter v0.123.0
16 - - gomod: github.com/open-telemetry/opentelemetry-collector-contrib/exporter/fileexporter v0.123.0
17 - - gomod: go.opentelemetry.io/collector/exporter/otlpexporter v0.123.0
15 + - gomod: go.opentelemetry.io/collector/exporter/debugexporter v0.124.0
16 + - gomod: github.com/open-telemetry/opentelemetry-collector-contrib/exporter/fileexporter v0.124.0
17 + - gomod: go.opentelemetry.io/collector/exporter/otlpexporter v0.124.0
18 - gomod: github.com/netdata/netdata/otel-collector/exporter/journaldexporter v0.0.0
19
20 replaces: