Unify behavior of the uploaders (#60534)

* Unify behavior of the uploaders

- Changes the behavior of uploaders to allow uploading pending
  summaries. It also makes all uploaders consistently fail to overwrite
  existing recordings and final summaries.
- Upgrades the gofakes3 library to support the If-None-Match headers.
- Copies the MemBuffer utility from the enterprise code.

* Fix a deadlock in gcsessions.Handler

The deadlock occurred when the upload passed a certain number of parts,
in which case the errors channel buffer got saturated and the object
deletion goroutines started piling up, ultimately leading to locking the
main loop that iterated through the object slice when the semaphore
channel got its buffer filled up too.

This deadlock prevented running the GCS uploader integration tests.

* review

* Another review round

- Move the pending session downloading to the caller
- Add a missing MemBuffer test
- Tweaks to the memory uploader (required by the Enterprise change)

* Build fix

* review

* Fix lint and integration tests

* Delete the pending summaries when uploading final, use versioning on s3

* Some comments

* license
This commit is contained in:
Bartosz Leper
2025-11-12 17:33:32 +00:00
committed by GitHub
parent 39dbce7268
commit 045d0b6ed0
19 changed files with 597 additions and 127 deletions
+1 -1
Submodule e updated: 51f6d75465...3386956073
+2 -2
View File
@@ -172,7 +172,7 @@ require (
github.com/jackc/pgx/v4 v4.18.3
github.com/jackc/pgx/v5 v5.7.6
github.com/jcmturner/gokrb5/v8 v8.4.4
github.com/johannesboyne/gofakes3 v0.0.0-20240217095638-c55a48f17be6
github.com/johannesboyne/gofakes3 v0.0.0-20250916175020-ebf3e50324d3
github.com/jonboulle/clockwork v0.5.0
github.com/joshlf/go-acl v0.0.0-20200411065538-eae00ae38531
github.com/json-iterator/go v1.1.12
@@ -560,7 +560,6 @@ require (
github.com/secure-systems-lab/go-securesystemslib v0.9.1 // indirect
github.com/segmentio/asm v1.2.1 // indirect
github.com/segmentio/ksuid v1.0.4 // indirect
github.com/shabbyrobe/gocovmerge v0.0.0-20230507112040-c3350d9342df // indirect
github.com/shibumi/go-pathspec v1.3.0 // indirect
github.com/shopspring/decimal v1.4.0 // indirect
github.com/siddontang/go v0.0.0-20180604090527-bdc77568d726 // indirect
@@ -611,6 +610,7 @@ require (
go.opentelemetry.io/contrib/detectors/gcp v1.38.0 // indirect
go.opentelemetry.io/otel/metric v1.38.0 // indirect
go.opentelemetry.io/otel/sdk/metric v1.38.0 // indirect
go.shabbyrobe.org/gocovmerge v0.0.0-20230507111327-fa4f82cfbf4d // indirect
go.uber.org/atomic v1.11.0 // indirect
go.uber.org/multierr v1.11.0 // indirect
go.uber.org/zap v1.27.0 // indirect
+9 -11
View File
@@ -786,7 +786,6 @@ github.com/asaskevich/govalidator v0.0.0-20230301143203-a9d515a09cc2 h1:DklsrG3d
github.com/asaskevich/govalidator v0.0.0-20230301143203-a9d515a09cc2/go.mod h1:WaHUgvxTVq04UNunO+XhnAqY/wQc+bxr74GqbsZ/Jqw=
github.com/atotto/clipboard v0.1.4 h1:EH0zSVneZPSuFR11BlR9YppQTVDbh5+16AmcJi4g1z4=
github.com/atotto/clipboard v0.1.4/go.mod h1:ZY9tmq7sm5xIbd9bOK4onWV4S6X0u6GY7Vn0Yu86PYI=
github.com/aws/aws-sdk-go v1.44.256/go.mod h1:aVsgQcEevwlmQ7qHE9I3h+dtQgpqhFB+i8Phjh7fkwI=
github.com/aws/aws-sdk-go v1.44.263/go.mod h1:aVsgQcEevwlmQ7qHE9I3h+dtQgpqhFB+i8Phjh7fkwI=
github.com/aws/aws-sdk-go v1.49.12/go.mod h1:LF8svs817+Nz+DmiMQKTO3ubZ/6IaTpq3TjupRn3Eqk=
github.com/aws/aws-sdk-go v1.55.8 h1:JRmEUbU52aJQZ2AjX4q4Wu7t4uZjOu71uyNmaWlUkJQ=
@@ -970,6 +969,8 @@ github.com/cespare/xxhash/v2 v2.1.1/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XL
github.com/cespare/xxhash/v2 v2.2.0/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs=
github.com/cespare/xxhash/v2 v2.3.0 h1:UL815xU9SqsFlibzuggzjXhog7bL6oX9BbNZnL2UFvs=
github.com/cespare/xxhash/v2 v2.3.0/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs=
github.com/cevatbarisyilmaz/ara v0.0.4 h1:SGH10hXpBJhhTlObuZzTuFn1rrdmjQImITXnZVPSodc=
github.com/cevatbarisyilmaz/ara v0.0.4/go.mod h1:BfFOxnUd6Mj6xmcvRxHN3Sr21Z1T3U2MYkYOmoQe4Ts=
github.com/chai2010/gettext-go v1.0.2 h1:1Lwwip6Q2QGsAdl/ZKPCwTe9fe0CjlUbqj5bFNSjIRk=
github.com/chai2010/gettext-go v1.0.2/go.mod h1:y+wnP2cHYaVj19NZhYKAwEMH2CI1gNHeQQ+5AjwawxA=
github.com/charlievieth/strcase v0.0.5 h1:gV4iXVyD6eI5KdfOV+/vIVCKXZwtCWOmDMcu7Uy00Rs=
@@ -1703,8 +1704,8 @@ github.com/jmhodges/clock v1.2.0 h1:eq4kys+NI0PLngzaHEe7AmPT90XMGIEySD1JfV1PDIs=
github.com/jmhodges/clock v1.2.0/go.mod h1:qKjhA7x7u/lQpPB1XAqX1b1lCI/w3/fNuYpI/ZjLynI=
github.com/jmoiron/sqlx v1.4.0 h1:1PLqN7S1UYp5t4SrVVnt4nUVNemrDAtxlulVe+Qgm3o=
github.com/jmoiron/sqlx v1.4.0/go.mod h1:ZrZ7UsYB/weZdl2Bxg6jCRO9c3YHl8r3ahlKmRT4JLY=
github.com/johannesboyne/gofakes3 v0.0.0-20240217095638-c55a48f17be6 h1:W8heH5NR7dfdB4FehSFI+DxjCbVKe9fPkPqKzCPJwnM=
github.com/johannesboyne/gofakes3 v0.0.0-20240217095638-c55a48f17be6/go.mod h1:AxgWC4DDX54O2WDoQO1Ceabtn6IbktjU/7bigor+66g=
github.com/johannesboyne/gofakes3 v0.0.0-20250916175020-ebf3e50324d3 h1:2713fQZ560HxoNVgfJH41GKzjMjIG+DW4hH6nYXfXW8=
github.com/johannesboyne/gofakes3 v0.0.0-20250916175020-ebf3e50324d3/go.mod h1:S4S9jGBVlLri0OeqrSSbCGG5vsI6he06UJyuz1WT1EE=
github.com/jonboulle/clockwork v0.5.0 h1:Hyh9A8u51kptdkR+cqRpT1EebBwTn1oK9YfGYbdFz6I=
github.com/jonboulle/clockwork v0.5.0/go.mod h1:3mZlmanh0g2NDKO5TWZVJAfofYk64M7XN3SzBPjZF60=
github.com/josharian/intern v1.0.0 h1:vlS4z54oSdjm0bgjRigI+G1HpF+tI+9rE5LLzOg8HmY=
@@ -2068,9 +2069,6 @@ github.com/segmentio/ksuid v1.0.4 h1:sBo2BdShXjmcugAMwjugoGUdUV0pcxY5mW4xKRn3v4c
github.com/segmentio/ksuid v1.0.4/go.mod h1:/XUiZBD3kVx5SmUOl55voK5yeAbBNNIed+2O73XgrPE=
github.com/sergi/go-diff v1.3.2-0.20230802210424-5b0b94c5c0d3 h1:n661drycOFuPLCN3Uc8sB6B/s6Z4t2xvBgU1htSHuq8=
github.com/sergi/go-diff v1.3.2-0.20230802210424-5b0b94c5c0d3/go.mod h1:A0bzQcvG0E7Rwjx0REVgAGH58e96+X0MeOfepqsbeW4=
github.com/shabbyrobe/gocovmerge v0.0.0-20190829150210-3e036491d500/go.mod h1:+njLrG5wSeoG4Ds61rFgEzKvenR2UHbjMoDHsczxly0=
github.com/shabbyrobe/gocovmerge v0.0.0-20230507112040-c3350d9342df h1:S77Pf5fIGMa7oSwp8SQPp7Hb4ZiI38K3RNBKD2LLeEM=
github.com/shabbyrobe/gocovmerge v0.0.0-20230507112040-c3350d9342df/go.mod h1:dcuzJZ83w/SqN9k4eQqwKYMgmKWzg/KzJAURBhRL1tc=
github.com/shibumi/go-pathspec v1.3.0 h1:QUyMZhFo0Md5B8zV8x2tesohbb5kfbpTi9rBnKh5dkI=
github.com/shibumi/go-pathspec v1.3.0/go.mod h1:Xutfslp817l2I1cZvgcfeMQJG5QnU2lh5tVaaMCl3jE=
github.com/shirou/gopsutil/v4 v4.25.10 h1:at8lk/5T1OgtuCp+AwrDofFRjnvosn0nkN2OLQ6g8tA=
@@ -2120,7 +2118,6 @@ github.com/sony/gobreaker v0.5.0/go.mod h1:ZKptC7FHNvhBz7dN2LGjPVBz2sZJmc0/PkyDJ
github.com/sourcegraph/conc v0.3.0 h1:OQTbbt6P72L20UqAkXXuLOj79LfEanQ+YQFNpLA9ySo=
github.com/sourcegraph/conc v0.3.0/go.mod h1:Sdozi7LEKbFPqYX2/J+iBAM6HpqSLTASQIKqDmF7Mt0=
github.com/spaolacci/murmur3 v0.0.0-20180118202830-f09979ecbc72/go.mod h1:JwIasOWyU6f++ZhiEuf87xNszmSA2myDM2Kzu9HwQUA=
github.com/spf13/afero v1.2.1/go.mod h1:9ZxEEn6pIJ8Rxe320qSDBk6AsU0r9pR7Q4OcevTdifk=
github.com/spf13/afero v1.3.3/go.mod h1:5KUK8ByomD5Ti5Artl0RtHeI5pTF7MIDuXL3yY520V4=
github.com/spf13/afero v1.6.0/go.mod h1:Ai8FlHk4v/PARR026UzYexafAt9roJ7LcLMAmO6Z93I=
github.com/spf13/afero v1.9.2/go.mod h1:iUV7ddyEEZPO5gA3zD4fJt6iStLlL+Lg4m2cihcDf8Y=
@@ -2276,7 +2273,8 @@ gitlab.com/gitlab-org/api/client-go v0.159.0 h1:ibKeribio/OCsrsUz7pkgIN4E7HWDyrw
gitlab.com/gitlab-org/api/client-go v0.159.0/go.mod h1:D0DHF7ILUfFo/JcoGMAEndiKMm8SiP/WjyJ4OfXxCKw=
go.einride.tech/aip v0.68.1 h1:16/AfSxcQISGN5z9C5lM+0mLYXihrHbQ1onvYTr93aQ=
go.einride.tech/aip v0.68.1/go.mod h1:XaFtaj4HuA3Zwk9xoBtTWgNubZ0ZZXv9BZJCkuKuWbg=
go.etcd.io/bbolt v1.3.5/go.mod h1:G5EMThwa9y8QZGBClrRx5EY+Yw9kAhnjy3bSjsnlVTQ=
go.etcd.io/bbolt v1.4.2 h1:IrUHp260R8c+zYx/Tm8QZr04CX+qWS5PGfPdevhdm1I=
go.etcd.io/bbolt v1.4.2/go.mod h1:Is8rSHO/b4f3XigBC0lL0+4FwAQv3HXEEIgFMuKHceM=
go.etcd.io/etcd/api/v3 v3.6.5 h1:pMMc42276sgR1j1raO/Qv3QI9Af/AuyQUW6CBAWuntA=
go.etcd.io/etcd/api/v3 v3.6.5/go.mod h1:ob0/oWA/UQQlT1BmaEkWQzI0sJ1M0Et0mMpaABxguOQ=
go.etcd.io/etcd/client/pkg/v3 v3.6.5 h1:Duz9fAzIZFhYWgRjp/FgNq2gO1jId9Yae/rLn3RrBP8=
@@ -2348,6 +2346,8 @@ go.opentelemetry.io/proto/otlp v0.15.0/go.mod h1:H7XAot3MsfNsj7EXtrA2q5xSNQ10UqI
go.opentelemetry.io/proto/otlp v0.19.0/go.mod h1:H7XAot3MsfNsj7EXtrA2q5xSNQ10UqI405h3+duxN4U=
go.opentelemetry.io/proto/otlp v1.9.0 h1:l706jCMITVouPOqEnii2fIAuO3IVGBRPV5ICjceRb/A=
go.opentelemetry.io/proto/otlp v1.9.0/go.mod h1:xE+Cx5E/eEHw+ISFkwPLwCZefwVjY+pqKg1qcK03+/4=
go.shabbyrobe.org/gocovmerge v0.0.0-20230507111327-fa4f82cfbf4d h1:Ns9kd1Rwzw7t0BR8XMphenji4SmIoNZPn8zhYmaVKP8=
go.shabbyrobe.org/gocovmerge v0.0.0-20230507111327-fa4f82cfbf4d/go.mod h1:92Uoe3l++MlthCm+koNi0tcUCX3anayogF0Pa/sp24k=
go.step.sm/crypto v0.70.0 h1:Q9Ft7N637mucyZcHZd1+0VVQJVwDCKqcb9CYcYi7cds=
go.step.sm/crypto v0.70.0/go.mod h1:pzfUhS5/ue7ev64PLlEgXvhx1opwbhFCjkvlhsxVds0=
go.uber.org/atomic v1.3.2/go.mod h1:gD2HeocX3+yG+ygLZcrzQJaqmWj9AIm7n08wl/qW/PE=
@@ -2460,7 +2460,6 @@ golang.org/x/mod v0.6.0-dev.0.20220419223038-86c51ed26bb4/go.mod h1:jJ57K6gSWd91
golang.org/x/mod v0.7.0/go.mod h1:iBbtSCu2XBx23ZKBPSOrRkjjQPZFPuis4dIYUhu/chs=
golang.org/x/mod v0.8.0/go.mod h1:iBbtSCu2XBx23ZKBPSOrRkjjQPZFPuis4dIYUhu/chs=
golang.org/x/mod v0.9.0/go.mod h1:iBbtSCu2XBx23ZKBPSOrRkjjQPZFPuis4dIYUhu/chs=
golang.org/x/mod v0.10.0/go.mod h1:iBbtSCu2XBx23ZKBPSOrRkjjQPZFPuis4dIYUhu/chs=
golang.org/x/mod v0.30.0 h1:fDEXFVZ/fmCKProc/yAXXUijritrDzahmwwefnjoPFk=
golang.org/x/mod v0.30.0/go.mod h1:lAsf5O2EvJeSFMiBxXDki7sCgAxEUcZHXoXMKT4GJKc=
golang.org/x/net v0.0.0-20180724234803-3673e40ba225/go.mod h1:mL1N/T3taQHkDXs73rZJwtUhF3w3ftmwwsq0BUmARs4=
@@ -2742,7 +2741,6 @@ golang.org/x/tools v0.0.0-20190621195816-6e04913cbbac/go.mod h1:/rFqwRUd4F7ZHNgw
golang.org/x/tools v0.0.0-20190628153133-6cdbf07be9d0/go.mod h1:/rFqwRUd4F7ZHNgwSSTFct+R/Kf4OFW1sUzUTQQTgfc=
golang.org/x/tools v0.0.0-20190816200558-6889da9d5479/go.mod h1:b+2E5dAYhXwXZwtnZ6UAqBI28+e2cm9otk0dWdXHAEo=
golang.org/x/tools v0.0.0-20190823170909-c4a336ef6a2f/go.mod h1:b+2E5dAYhXwXZwtnZ6UAqBI28+e2cm9otk0dWdXHAEo=
golang.org/x/tools v0.0.0-20190829051458-42f498d34c4d/go.mod h1:b+2E5dAYhXwXZwtnZ6UAqBI28+e2cm9otk0dWdXHAEo=
golang.org/x/tools v0.0.0-20190911174233-4f2ddba30aff/go.mod h1:b+2E5dAYhXwXZwtnZ6UAqBI28+e2cm9otk0dWdXHAEo=
golang.org/x/tools v0.0.0-20190927191325-030b2cf1153e/go.mod h1:b+2E5dAYhXwXZwtnZ6UAqBI28+e2cm9otk0dWdXHAEo=
golang.org/x/tools v0.0.0-20191012152004-8de300cfc20a/go.mod h1:b+2E5dAYhXwXZwtnZ6UAqBI28+e2cm9otk0dWdXHAEo=
@@ -2796,7 +2794,6 @@ golang.org/x/tools v0.1.12/go.mod h1:hNGJHUnrk76NpqgfD5Aqm5Crs+Hm0VOH/i9J2+nxYbc
golang.org/x/tools v0.3.0/go.mod h1:/rWhSS2+zyEVwoJf8YAX6L2f0ntZ7Kn/mGgAWcipA5k=
golang.org/x/tools v0.6.0/go.mod h1:Xwgl3UAJ/d3gWutnCtw505GrjyAbvKui8lOU390QaIU=
golang.org/x/tools v0.7.0/go.mod h1:4pg6aUX35JBAogB10C9AtvVL+qowtN4pT3CGSQex14s=
golang.org/x/tools v0.8.0/go.mod h1:JxBZ99ISMI5ViVkT1tr6tdNmXeTrcpVSD3vZ1RsRdN4=
golang.org/x/tools v0.38.0 h1:Hx2Xv8hISq8Lm16jvBZ2VQf+RLmbd7wVUsALibYI/IQ=
golang.org/x/tools v0.38.0/go.mod h1:yEsQ/d/YK8cjh0L6rZlY8tgtlKiBNTL14pGDJPJpYQs=
golang.org/x/tools/go/expect v0.1.0-deprecated h1:jY2C5HGYR5lqex3gEniOQL0r7Dq5+VGVgY1nudX5lXY=
@@ -3122,6 +3119,7 @@ gopkg.in/ini.v1 v1.67.0 h1:Dgnx+6+nfE+IfzjUEISNeydPJh9AXNNsWbGP9KzCsOA=
gopkg.in/ini.v1 v1.67.0/go.mod h1:pNLf8WUiyNEtQjuu5G5vTm06TEv9tsIgeAvK8hOrP4k=
gopkg.in/mail.v2 v2.3.1 h1:WYFn/oANrAGP2C0dcV6/pbkPzv8yGzqTjPmTeO7qoXk=
gopkg.in/mail.v2 v2.3.1/go.mod h1:htwXN1Qh09vZJ1NVKxQqHPBaCBbzKhp5GzuJEA4VJWw=
gopkg.in/mgo.v2 v2.0.0-20180705113604-9856a29383ce h1:xcEWjVhvbDy+nHP67nPDDpbYrY+ILlfndk4bRioVHaU=
gopkg.in/mgo.v2 v2.0.0-20180705113604-9856a29383ce/go.mod h1:yeKp02qBN3iKW1OzL3MGk2IdtZzaj7SFntXj72NppTA=
gopkg.in/natefinch/lumberjack.v2 v2.0.0/go.mod h1:l0ndWWf7gzL7RNwBG7wST/UCcT4T24xpD6X8LsfU/+k=
gopkg.in/natefinch/lumberjack.v2 v2.2.1 h1:bBRl1b0OH9s/DuPhuXpNl+VtCaJXFZ5/uEFST95x9zc=
+8 -2
View File
@@ -5116,10 +5116,16 @@ readLoop:
require.True(t, hasLeave, "session leave event not found")
require.True(t, hasEnd, "session end event not found")
// ensure session upload directory is empty
// ensure session upload directory is empty, apart from the pending directory
fi, err := os.ReadDir(sessionsDir)
require.NoError(t, err)
require.Empty(t, fi)
require.Len(t, fi, 1)
assert.Equal(t, "pending", fi[0].Name())
assert.True(t, fi[0].IsDir())
fi, err = os.ReadDir(filepath.Join(sessionsDir, "pending"))
require.NoError(t, err)
assert.Empty(t, fi)
}
// testPAM checks that Teleport PAM integration works correctly. In this case
+4 -4
View File
@@ -1688,8 +1688,8 @@ github.com/jmhodges/clock v1.2.0 h1:eq4kys+NI0PLngzaHEe7AmPT90XMGIEySD1JfV1PDIs=
github.com/jmhodges/clock v1.2.0/go.mod h1:qKjhA7x7u/lQpPB1XAqX1b1lCI/w3/fNuYpI/ZjLynI=
github.com/jmoiron/sqlx v1.4.0 h1:1PLqN7S1UYp5t4SrVVnt4nUVNemrDAtxlulVe+Qgm3o=
github.com/jmoiron/sqlx v1.4.0/go.mod h1:ZrZ7UsYB/weZdl2Bxg6jCRO9c3YHl8r3ahlKmRT4JLY=
github.com/johannesboyne/gofakes3 v0.0.0-20240217095638-c55a48f17be6 h1:W8heH5NR7dfdB4FehSFI+DxjCbVKe9fPkPqKzCPJwnM=
github.com/johannesboyne/gofakes3 v0.0.0-20240217095638-c55a48f17be6/go.mod h1:AxgWC4DDX54O2WDoQO1Ceabtn6IbktjU/7bigor+66g=
github.com/johannesboyne/gofakes3 v0.0.0-20250916175020-ebf3e50324d3 h1:2713fQZ560HxoNVgfJH41GKzjMjIG+DW4hH6nYXfXW8=
github.com/johannesboyne/gofakes3 v0.0.0-20250916175020-ebf3e50324d3/go.mod h1:S4S9jGBVlLri0OeqrSSbCGG5vsI6he06UJyuz1WT1EE=
github.com/jonboulle/clockwork v0.5.0 h1:Hyh9A8u51kptdkR+cqRpT1EebBwTn1oK9YfGYbdFz6I=
github.com/jonboulle/clockwork v0.5.0/go.mod h1:3mZlmanh0g2NDKO5TWZVJAfofYk64M7XN3SzBPjZF60=
github.com/josharian/intern v1.0.0 h1:vlS4z54oSdjm0bgjRigI+G1HpF+tI+9rE5LLzOg8HmY=
@@ -2009,8 +2009,6 @@ github.com/segmentio/ksuid v1.0.4 h1:sBo2BdShXjmcugAMwjugoGUdUV0pcxY5mW4xKRn3v4c
github.com/segmentio/ksuid v1.0.4/go.mod h1:/XUiZBD3kVx5SmUOl55voK5yeAbBNNIed+2O73XgrPE=
github.com/sergi/go-diff v1.3.2-0.20230802210424-5b0b94c5c0d3 h1:n661drycOFuPLCN3Uc8sB6B/s6Z4t2xvBgU1htSHuq8=
github.com/sergi/go-diff v1.3.2-0.20230802210424-5b0b94c5c0d3/go.mod h1:A0bzQcvG0E7Rwjx0REVgAGH58e96+X0MeOfepqsbeW4=
github.com/shabbyrobe/gocovmerge v0.0.0-20230507112040-c3350d9342df h1:S77Pf5fIGMa7oSwp8SQPp7Hb4ZiI38K3RNBKD2LLeEM=
github.com/shabbyrobe/gocovmerge v0.0.0-20230507112040-c3350d9342df/go.mod h1:dcuzJZ83w/SqN9k4eQqwKYMgmKWzg/KzJAURBhRL1tc=
github.com/shibumi/go-pathspec v1.3.0 h1:QUyMZhFo0Md5B8zV8x2tesohbb5kfbpTi9rBnKh5dkI=
github.com/shibumi/go-pathspec v1.3.0/go.mod h1:Xutfslp817l2I1cZvgcfeMQJG5QnU2lh5tVaaMCl3jE=
github.com/shirou/gopsutil v3.21.11+incompatible h1:+1+c1VGhc88SSonWP6foOcLhvnKlUeu/erjjvaPEYiI=
@@ -2285,6 +2283,8 @@ go.opentelemetry.io/proto/otlp v0.15.0/go.mod h1:H7XAot3MsfNsj7EXtrA2q5xSNQ10UqI
go.opentelemetry.io/proto/otlp v0.19.0/go.mod h1:H7XAot3MsfNsj7EXtrA2q5xSNQ10UqI405h3+duxN4U=
go.opentelemetry.io/proto/otlp v1.9.0 h1:l706jCMITVouPOqEnii2fIAuO3IVGBRPV5ICjceRb/A=
go.opentelemetry.io/proto/otlp v1.9.0/go.mod h1:xE+Cx5E/eEHw+ISFkwPLwCZefwVjY+pqKg1qcK03+/4=
go.shabbyrobe.org/gocovmerge v0.0.0-20230507111327-fa4f82cfbf4d h1:Ns9kd1Rwzw7t0BR8XMphenji4SmIoNZPn8zhYmaVKP8=
go.shabbyrobe.org/gocovmerge v0.0.0-20230507111327-fa4f82cfbf4d/go.mod h1:92Uoe3l++MlthCm+koNi0tcUCX3anayogF0Pa/sp24k=
go.step.sm/crypto v0.70.0 h1:Q9Ft7N637mucyZcHZd1+0VVQJVwDCKqcb9CYcYi7cds=
go.step.sm/crypto v0.70.0/go.mod h1:pzfUhS5/ue7ev64PLlEgXvhx1opwbhFCjkvlhsxVds0=
go.uber.org/atomic v1.3.2/go.mod h1:gD2HeocX3+yG+ygLZcrzQJaqmWj9AIm7n08wl/qW/PE=
+4 -4
View File
@@ -1742,8 +1742,8 @@ github.com/jmhodges/clock v1.2.0 h1:eq4kys+NI0PLngzaHEe7AmPT90XMGIEySD1JfV1PDIs=
github.com/jmhodges/clock v1.2.0/go.mod h1:qKjhA7x7u/lQpPB1XAqX1b1lCI/w3/fNuYpI/ZjLynI=
github.com/jmoiron/sqlx v1.4.0 h1:1PLqN7S1UYp5t4SrVVnt4nUVNemrDAtxlulVe+Qgm3o=
github.com/jmoiron/sqlx v1.4.0/go.mod h1:ZrZ7UsYB/weZdl2Bxg6jCRO9c3YHl8r3ahlKmRT4JLY=
github.com/johannesboyne/gofakes3 v0.0.0-20240217095638-c55a48f17be6 h1:W8heH5NR7dfdB4FehSFI+DxjCbVKe9fPkPqKzCPJwnM=
github.com/johannesboyne/gofakes3 v0.0.0-20240217095638-c55a48f17be6/go.mod h1:AxgWC4DDX54O2WDoQO1Ceabtn6IbktjU/7bigor+66g=
github.com/johannesboyne/gofakes3 v0.0.0-20250916175020-ebf3e50324d3 h1:2713fQZ560HxoNVgfJH41GKzjMjIG+DW4hH6nYXfXW8=
github.com/johannesboyne/gofakes3 v0.0.0-20250916175020-ebf3e50324d3/go.mod h1:S4S9jGBVlLri0OeqrSSbCGG5vsI6he06UJyuz1WT1EE=
github.com/jonboulle/clockwork v0.5.0 h1:Hyh9A8u51kptdkR+cqRpT1EebBwTn1oK9YfGYbdFz6I=
github.com/jonboulle/clockwork v0.5.0/go.mod h1:3mZlmanh0g2NDKO5TWZVJAfofYk64M7XN3SzBPjZF60=
github.com/josharian/intern v1.0.0 h1:vlS4z54oSdjm0bgjRigI+G1HpF+tI+9rE5LLzOg8HmY=
@@ -2090,8 +2090,6 @@ github.com/sergi/go-diff v1.1.0/go.mod h1:STckp+ISIX8hZLjrqAeVduY0gWCT9IjLuqbuNX
github.com/sergi/go-diff v1.2.0/go.mod h1:STckp+ISIX8hZLjrqAeVduY0gWCT9IjLuqbuNXdaHfM=
github.com/sergi/go-diff v1.3.2-0.20230802210424-5b0b94c5c0d3 h1:n661drycOFuPLCN3Uc8sB6B/s6Z4t2xvBgU1htSHuq8=
github.com/sergi/go-diff v1.3.2-0.20230802210424-5b0b94c5c0d3/go.mod h1:A0bzQcvG0E7Rwjx0REVgAGH58e96+X0MeOfepqsbeW4=
github.com/shabbyrobe/gocovmerge v0.0.0-20230507112040-c3350d9342df h1:S77Pf5fIGMa7oSwp8SQPp7Hb4ZiI38K3RNBKD2LLeEM=
github.com/shabbyrobe/gocovmerge v0.0.0-20230507112040-c3350d9342df/go.mod h1:dcuzJZ83w/SqN9k4eQqwKYMgmKWzg/KzJAURBhRL1tc=
github.com/shibumi/go-pathspec v1.3.0 h1:QUyMZhFo0Md5B8zV8x2tesohbb5kfbpTi9rBnKh5dkI=
github.com/shibumi/go-pathspec v1.3.0/go.mod h1:Xutfslp817l2I1cZvgcfeMQJG5QnU2lh5tVaaMCl3jE=
github.com/shirou/gopsutil v3.21.11+incompatible h1:+1+c1VGhc88SSonWP6foOcLhvnKlUeu/erjjvaPEYiI=
@@ -2373,6 +2371,8 @@ go.opentelemetry.io/proto/otlp v0.15.0/go.mod h1:H7XAot3MsfNsj7EXtrA2q5xSNQ10UqI
go.opentelemetry.io/proto/otlp v0.19.0/go.mod h1:H7XAot3MsfNsj7EXtrA2q5xSNQ10UqI405h3+duxN4U=
go.opentelemetry.io/proto/otlp v1.9.0 h1:l706jCMITVouPOqEnii2fIAuO3IVGBRPV5ICjceRb/A=
go.opentelemetry.io/proto/otlp v1.9.0/go.mod h1:xE+Cx5E/eEHw+ISFkwPLwCZefwVjY+pqKg1qcK03+/4=
go.shabbyrobe.org/gocovmerge v0.0.0-20230507111327-fa4f82cfbf4d h1:Ns9kd1Rwzw7t0BR8XMphenji4SmIoNZPn8zhYmaVKP8=
go.shabbyrobe.org/gocovmerge v0.0.0-20230507111327-fa4f82cfbf4d/go.mod h1:92Uoe3l++MlthCm+koNi0tcUCX3anayogF0Pa/sp24k=
go.step.sm/crypto v0.70.0 h1:Q9Ft7N637mucyZcHZd1+0VVQJVwDCKqcb9CYcYi7cds=
go.step.sm/crypto v0.70.0/go.mod h1:pzfUhS5/ue7ev64PLlEgXvhx1opwbhFCjkvlhsxVds0=
go.uber.org/atomic v1.3.2/go.mod h1:gD2HeocX3+yG+ygLZcrzQJaqmWj9AIm7n08wl/qW/PE=
+77 -15
View File
@@ -269,6 +269,12 @@ func (h *Handler) sessionBlob(sessionID session.ID) *blockblob.Client {
return h.session.NewBlockBlobClient(sessionName(sessionID))
}
// pendingSummaryBlob returns a BlockBlobClient for the blob of the a pending
// session summary.
func (h *Handler) pendingSummaryBlob(sessionID session.ID) *blockblob.Client {
return h.inprogress.NewBlockBlobClient(summaryName(sessionID))
}
// summaryBlob returns a BlockBlobClient for the blob of the session summary.
func (h *Handler) summaryBlob(sessionID session.ID) *blockblob.Client {
return h.session.NewBlockBlobClient(summaryName(sessionID))
@@ -298,31 +304,74 @@ func (h *Handler) partBlob(upload events.StreamUpload, partNumber int64) *blockb
// Upload implements [events.UploadHandler] and uploads a session recording.
func (h *Handler) Upload(ctx context.Context, sessionID session.ID, reader io.Reader) (string, error) {
return h.uploadBlob(ctx, sessionID, h.sessionBlob(sessionID), reader)
path, err := h.uploadBlob(ctx, sessionID, h.sessionBlob(sessionID), reader)
return path, trace.Wrap(err)
}
// UploadSummary implements [events.UploadHandler] and uploads a session
// summary.
// UploadPendingSummary implements [events.UploadHandler] and uploads a pending
// session summary. This function can be called multiple times for a given
// sessionID to update the state.
func (h *Handler) UploadPendingSummary(ctx context.Context, sessionID session.ID, reader io.Reader) (string, error) {
path, err := h.uploadBlob(ctx, sessionID, h.pendingSummaryBlob(sessionID), reader, withOverwrite())
return path, trace.Wrap(err)
}
// UploadSummary implements [events.UploadHandler] and uploads a final version
// of session summary and deletes the pending one. This function can be called
// only once for a given sessionID; subsequent calls will return an error.
func (h *Handler) UploadSummary(ctx context.Context, sessionID session.ID, reader io.Reader) (string, error) {
return h.uploadBlob(ctx, sessionID, h.summaryBlob(sessionID), reader)
path, err := h.uploadBlob(ctx, sessionID, h.summaryBlob(sessionID), reader)
if err != nil {
return "", trace.Wrap(err)
}
_, err = cErr(h.pendingSummaryBlob(sessionID).Delete(ctx, nil))
if err != nil && !trace.IsNotFound(err) {
return "", trace.Wrap(err)
}
return path, nil
}
// UploadMetadata implements [events.UploadHandler] and uploads the session
// metadata.
func (h *Handler) UploadMetadata(ctx context.Context, sessionID session.ID, reader io.Reader) (string, error) {
return h.uploadBlob(ctx, sessionID, h.metadataBlob(sessionID), reader)
path, err := h.uploadBlob(ctx, sessionID, h.metadataBlob(sessionID), reader)
return path, trace.Wrap(err)
}
// UploadThumbnail implements [events.UploadHandler] and uploads the session
// thumbnail.
func (h *Handler) UploadThumbnail(ctx context.Context, sessionID session.ID, reader io.Reader) (string, error) {
return h.uploadBlob(ctx, sessionID, h.thumbnailBlob(sessionID), reader)
path, err := h.uploadBlob(ctx, sessionID, h.thumbnailBlob(sessionID), reader)
return path, trace.Wrap(err)
}
func (h *Handler) uploadBlob(ctx context.Context, sessionID session.ID, blob *blockblob.Client, reader io.Reader) (string, error) {
if _, err := cErr(blob.UploadStream(ctx, reader, &blockblob.UploadStreamOptions{
AccessConditions: &blobDoesNotExist,
})); err != nil {
type blobUploadConfig struct {
overwrite bool
}
type blobUploadOption func(*blobUploadConfig)
func withOverwrite() blobUploadOption {
return func(cfg *blobUploadConfig) {
cfg.overwrite = true
}
}
func (h *Handler) uploadBlob(
ctx context.Context, sessionID session.ID, blob *blockblob.Client, reader io.Reader, opts ...blobUploadOption,
) (string, error) {
cfg := blobUploadConfig{}
for _, opt := range opts {
opt(&cfg)
}
usOpts := blockblob.UploadStreamOptions{}
if !cfg.overwrite {
usOpts.AccessConditions = &blobDoesNotExist
}
if _, err := cErr(blob.UploadStream(ctx, reader, &usOpts)); err != nil {
return "", trace.Wrap(err)
}
h.log.DebugContext(ctx, "Blob uploaded.", fieldSessionID, sessionID)
@@ -332,22 +381,35 @@ func (h *Handler) uploadBlob(ctx context.Context, sessionID session.ID, blob *bl
// Download implements [events.UploadHandler] and downloads a session recording.
func (h *Handler) Download(ctx context.Context, sessionID session.ID, writer events.RandomAccessWriter) error {
return h.downloadBlob(ctx, sessionID, h.sessionBlob(sessionID), writer)
return trace.Wrap(h.downloadBlob(ctx, sessionID, h.sessionBlob(sessionID), writer))
}
// DownloadSummary implements [events.UploadHandler] and downloads a session summary.
// DownloadSummary implements [events.UploadHandler] and downloads a final
// session summary.
func (h *Handler) DownloadSummary(ctx context.Context, sessionID session.ID, writer events.RandomAccessWriter) error {
return h.downloadBlob(ctx, sessionID, h.summaryBlob(sessionID), writer)
// Happy path: the final summary exists.
err := h.downloadBlob(ctx, sessionID, h.summaryBlob(sessionID), writer)
if trace.IsNotFound(err) {
// Final summary doesn't exist, try the pending one.
err = h.downloadBlob(ctx, sessionID, h.pendingSummaryBlob(sessionID), writer)
if trace.IsNotFound(err) {
// One more check for the final summary to prevent a race condition where
// the final one got created and the pending one got removed between the
// two checks above.
err = h.downloadBlob(ctx, sessionID, h.summaryBlob(sessionID), writer)
}
}
return trace.Wrap(err)
}
// DownloadMetadata implements [events.UploadHandler] and downloads a session's metadata.
func (h *Handler) DownloadMetadata(ctx context.Context, sessionID session.ID, writer events.RandomAccessWriter) error {
return h.downloadBlob(ctx, sessionID, h.metadataBlob(sessionID), writer)
return trace.Wrap(h.downloadBlob(ctx, sessionID, h.metadataBlob(sessionID), writer))
}
// DownloadThumbnail implements [events.UploadHandler] and downloads a session's thumbnail.
func (h *Handler) DownloadThumbnail(ctx context.Context, sessionID session.ID, writer events.RandomAccessWriter) error {
return h.downloadBlob(ctx, sessionID, h.thumbnailBlob(sessionID), writer)
return trace.Wrap(h.downloadBlob(ctx, sessionID, h.thumbnailBlob(sessionID), writer))
}
func (h *Handler) downloadBlob(ctx context.Context, sessionID session.ID, blob *blockblob.Client, writer events.RandomAccessWriter) error {
+36 -14
View File
@@ -49,12 +49,13 @@ type MemoryUploaderConfig struct {
// upload
func NewMemoryUploader(cfg ...MemoryUploaderConfig) *MemoryUploader {
up := &MemoryUploader{
mtx: &sync.RWMutex{},
uploads: make(map[string]*MemoryUpload),
sessions: make(map[session.ID][]byte),
summaries: make(map[session.ID][]byte),
metadata: make(map[session.ID][]byte),
thumbnails: make(map[session.ID][]byte),
mtx: &sync.RWMutex{},
uploads: make(map[string]*MemoryUpload),
sessions: make(map[session.ID][]byte),
summaries: make(map[session.ID][]byte),
pendingSummaries: make(map[session.ID][]byte),
metadata: make(map[session.ID][]byte),
thumbnails: make(map[session.ID][]byte),
}
if len(cfg) != 0 {
up.cfg = cfg[0]
@@ -66,12 +67,13 @@ func NewMemoryUploader(cfg ...MemoryUploaderConfig) *MemoryUploader {
type MemoryUploader struct {
cfg MemoryUploaderConfig
mtx *sync.RWMutex
uploads map[string]*MemoryUpload
sessions map[session.ID][]byte
summaries map[session.ID][]byte
metadata map[session.ID][]byte
thumbnails map[session.ID][]byte
mtx *sync.RWMutex
uploads map[string]*MemoryUpload
sessions map[session.ID][]byte
summaries map[session.ID][]byte
pendingSummaries map[session.ID][]byte
metadata map[session.ID][]byte
thumbnails map[session.ID][]byte
// Clock is an optional [clockwork.Clock] to determine the time to associate
// with uploads and parts.
@@ -115,6 +117,7 @@ func (m *MemoryUploader) Reset() {
m.uploads = make(map[string]*MemoryUpload)
m.sessions = make(map[session.ID][]byte)
m.summaries = make(map[session.ID][]byte)
m.pendingSummaries = make(map[session.ID][]byte)
m.metadata = make(map[session.ID][]byte)
m.thumbnails = make(map[session.ID][]byte)
}
@@ -285,8 +288,23 @@ func (m *MemoryUploader) Upload(ctx context.Context, sessionID session.ID, readC
return string(sessionID), nil
}
// UploadSummary uploads session summary and returns URL with uploaded file in
// case of success.
// UploadPendingSummary uploads a pending session summary. This function can be
// called multiple times for a given sessionID to update the state.
func (m *MemoryUploader) UploadPendingSummary(ctx context.Context, sessionID session.ID, readCloser io.Reader) (string, error) {
m.mtx.Lock()
defer m.mtx.Unlock()
// Pending summary can be overridden, so skip checking for existing one.
data, err := io.ReadAll(readCloser)
if err != nil {
return "", trace.ConvertSystemError(err)
}
m.pendingSummaries[sessionID] = data
return string(sessionID), nil
}
// UploadSummary uploads final version of session summary and returns URL with
// uploaded file in case of success. This function can be called only once for
// a given sessionID; subsequent calls will return an error.
func (m *MemoryUploader) UploadSummary(ctx context.Context, sessionID session.ID, readCloser io.Reader) (string, error) {
m.mtx.Lock()
defer m.mtx.Unlock()
@@ -299,6 +317,7 @@ func (m *MemoryUploader) UploadSummary(ctx context.Context, sessionID session.ID
return "", trace.ConvertSystemError(err)
}
m.summaries[sessionID] = data
delete(m.pendingSummaries, sessionID)
return string(sessionID), nil
}
@@ -357,6 +376,9 @@ func (m *MemoryUploader) DownloadSummary(ctx context.Context, sessionID session.
defer m.mtx.RUnlock()
data, ok := m.summaries[sessionID]
if !ok {
data, ok = m.pendingSummaries[sessionID]
}
if !ok {
return trace.NotFound("summary %q is not found", sessionID)
}
+76 -9
View File
@@ -91,6 +91,12 @@ func NewHandler(cfg Config) (*Handler, error) {
Config: cfg,
fileRecorder: NewPlainFileRecorder(logger, cfg.OpenFile),
}
err := os.MkdirAll(h.pendingSummariesPath(), teleport.PrivateDirMode)
if err != nil {
return nil, trace.ConvertSystemError(err)
}
return h, nil
}
@@ -112,22 +118,34 @@ func (l *Handler) Close() error {
// Download reads a session recording from a local directory.
func (l *Handler) Download(ctx context.Context, sessionID session.ID, writer events.RandomAccessWriter) error {
return downloadFile(l.recordingPath(sessionID), writer)
return trace.Wrap(downloadFile(l.recordingPath(sessionID), writer))
}
// DownloadSummary reads a session summary from a local directory.
func (l *Handler) DownloadSummary(ctx context.Context, sessionID session.ID, writer events.RandomAccessWriter) error {
return downloadFile(l.summaryPath(sessionID), writer)
// Happy path: the final summary exists.
err := downloadFile(l.summaryPath(sessionID), writer)
if trace.IsNotFound(err) {
// Final summary doesn't exist, try the pending one.
err = downloadFile(l.pendingSummaryPath(sessionID), writer)
if trace.IsNotFound(err) {
// One more check for the final summary to prevent a race condition where
// the final one got created and the pending one got removed between the
// two checks above.
err = downloadFile(l.summaryPath(sessionID), writer)
}
}
return trace.Wrap(err)
}
// DownloadMetadata reads session metadata from a local directory.
func (l *Handler) DownloadMetadata(ctx context.Context, sessionID session.ID, writer events.RandomAccessWriter) error {
return downloadFile(l.metadataPath(sessionID), writer)
return trace.Wrap(downloadFile(l.metadataPath(sessionID), writer))
}
// DownloadThumbnail reads a session thumbnail from a local directory.
func (l *Handler) DownloadThumbnail(ctx context.Context, sessionID session.ID, writer events.RandomAccessWriter) error {
return downloadFile(l.thumbnailPath(sessionID), writer)
return trace.Wrap(downloadFile(l.thumbnailPath(sessionID), writer))
}
func downloadFile(path string, writer events.RandomAccessWriter) error {
@@ -145,12 +163,32 @@ func downloadFile(path string, writer events.RandomAccessWriter) error {
// Upload writes a session recording to a local directory.
func (l *Handler) Upload(ctx context.Context, sessionID session.ID, reader io.Reader) (string, error) {
return uploadFile(l.recordingPath(sessionID), reader)
s, err := uploadFile(l.recordingPath(sessionID), reader)
return s, trace.Wrap(err)
}
// UploadSummary writes a session summary to a local directory.
// UploadPendingSummary writes a pending session summary to a local directory.
// This function can be called multiple times for a given sessionID to update
// the state.
func (l *Handler) UploadPendingSummary(ctx context.Context, sessionID session.ID, reader io.Reader) (string, error) {
return uploadFile(l.pendingSummaryPath(sessionID), reader, withOverwrite())
}
// UploadSummary writes a final version of session summary and removes the
// pending one. This function can be called only once for a given sessionID;
// subsequent calls will return an error.
func (l *Handler) UploadSummary(ctx context.Context, sessionID session.ID, reader io.Reader) (string, error) {
return uploadFile(l.summaryPath(sessionID), reader)
name, err := uploadFile(l.summaryPath(sessionID), reader)
if err != nil {
return "", trace.Wrap(err)
}
err = os.Remove(l.pendingSummaryPath(sessionID))
if err != nil && !trace.IsNotFound(err) {
return "", trace.Wrap(err)
}
return name, nil
}
// UploadMetadata writes session metadata to a local directory.
@@ -163,8 +201,29 @@ func (l *Handler) UploadThumbnail(ctx context.Context, sessionID session.ID, rea
return uploadFile(l.thumbnailPath(sessionID), reader)
}
func uploadFile(path string, reader io.Reader) (string, error) {
f, err := os.Create(path)
type fileUploadConfig struct {
overwrite bool
}
type fileUploadOption func(*fileUploadConfig)
func withOverwrite() fileUploadOption {
return func(cfg *fileUploadConfig) {
cfg.overwrite = true
}
}
func uploadFile(path string, reader io.Reader, opts ...fileUploadOption) (string, error) {
cfg := fileUploadConfig{}
for _, opt := range opts {
opt(&cfg)
}
flags := os.O_RDWR | os.O_CREATE | os.O_TRUNC
if !cfg.overwrite {
flags |= os.O_EXCL
}
f, err := os.OpenFile(path, flags, 0666)
if err != nil {
return "", trace.ConvertSystemError(err)
}
@@ -179,6 +238,14 @@ func (l *Handler) recordingPath(sessionID session.ID) string {
return filepath.Join(l.Directory, string(sessionID)+tarExt)
}
func (l *Handler) pendingSummariesPath() string {
return filepath.Join(l.Directory, "pending")
}
func (l *Handler) pendingSummaryPath(sessionID session.ID) string {
return filepath.Join(l.pendingSummariesPath(), string(sessionID)+summaryExt)
}
func (l *Handler) summaryPath(sessionID session.ID) string {
return filepath.Join(l.Directory, string(sessionID)+summaryExt)
}
+81 -22
View File
@@ -88,6 +88,8 @@ const (
kmsKeyName = "keyName"
// pathPropertyKey
pathPropertyKey = "path"
// pendingPrefix is a prefix of the pending summaries path
pendingPrefix = "pending"
)
// Config is handler configuration
@@ -238,43 +240,81 @@ func (h *Handler) Close() error {
// Upload reads the content of a session recording from a reader and uploads it
// to a GCS bucket. If successful, it returns URL of the uploaded object.
func (h *Handler) Upload(ctx context.Context, sessionID session.ID, reader io.Reader) (string, error) {
return h.uploadFile(ctx, h.recordingPath(sessionID), reader)
path, err := h.uploadFile(ctx, h.recordingPath(sessionID), reader)
return path, trace.Wrap(err)
}
// UploadSummary reads the content of a session summary from a reader and
// uploads it to a GCS bucket. If successful, it returns URL of the uploaded
// object.
// UploadPendingSummary reads the content of a pending session summary from a
// reader and uploads it to a GCS bucket. If successful, it returns URL of the
// uploaded object. This function can be called multiple times for a given
// sessionID to update the state.
func (h *Handler) UploadPendingSummary(ctx context.Context, sessionID session.ID, reader io.Reader) (string, error) {
path, err := h.uploadFile(ctx, h.pendingSummaryPath(sessionID), reader, withOverwrite())
return path, trace.Wrap(err)
}
// UploadSummary reads the content of a final version of session summary from a
// reader and uploads it to a GCS bucket. The pending version, if any, is
// removed. If successful, it returns URL of the uploaded object. This function
// can be called only once for a given sessionID; subsequent calls will return
// an error.
func (h *Handler) UploadSummary(ctx context.Context, sessionID session.ID, reader io.Reader) (string, error) {
return h.uploadFile(ctx, h.summaryPath(sessionID), reader)
path, err := h.uploadFile(ctx, h.summaryPath(sessionID), reader)
if err != nil {
return "", trace.Wrap(err)
}
pendingObj := h.gcsClient.Bucket(h.Config.Bucket).Object(h.pendingSummaryPath(sessionID))
err = convertGCSError(pendingObj.Delete(ctx))
if err != nil && !trace.IsNotFound(err) {
return "", trace.Wrap(err)
}
return path, nil
}
// UploadMetadata reads the session metadata from a reader and uploads it to a GCS
// bucket. If successful, it returns URL of the uploaded object.
func (h *Handler) UploadMetadata(ctx context.Context, sessionID session.ID, reader io.Reader) (string, error) {
return h.uploadFile(ctx, h.metadataPath(sessionID), reader)
path, err := h.uploadFile(ctx, h.metadataPath(sessionID), reader)
return path, trace.Wrap(err)
}
// UploadThumbnail reads the session thumbnail from a reader and uploads it to a GCS
// bucket. If successful, it returns URL of the uploaded object.
func (h *Handler) UploadThumbnail(ctx context.Context, sessionID session.ID, reader io.Reader) (string, error) {
return h.uploadFile(ctx, h.thumbnailPath(sessionID), reader)
path, err := h.uploadFile(ctx, h.thumbnailPath(sessionID), reader)
return path, trace.Wrap(err)
}
func (h *Handler) uploadFile(ctx context.Context, path string, reader io.Reader) (string, error) {
h.logger.DebugContext(ctx, "Uploading object to GCS", "path", path)
type fileUploadConfig struct {
overwrite bool
}
// Make sure we don't overwrite an existing recording.
_, err := h.gcsClient.Bucket(h.Config.Bucket).Object(path).Attrs(ctx)
if !errors.Is(err, storage.ErrObjectNotExist) {
if err != nil {
return "", convertGCSError(err)
}
return "", trace.AlreadyExists("file %q already exists in GCS", path)
type fileUploadOption func(*fileUploadConfig)
func withOverwrite() fileUploadOption {
return func(cfg *fileUploadConfig) {
cfg.overwrite = true
}
}
func (h *Handler) uploadFile(ctx context.Context, path string, reader io.Reader, opts ...fileUploadOption) (string, error) {
cfg := fileUploadConfig{}
for _, opt := range opts {
opt(&cfg)
}
writer := h.gcsClient.Bucket(h.Config.Bucket).Object(path).NewWriter(ctx)
h.logger.DebugContext(ctx, "Uploading object to GCS", "path", path)
obj := h.gcsClient.Bucket(h.Config.Bucket).Object(path)
if !cfg.overwrite {
// Make sure we don't overwrite an existing recording.
obj = obj.If(storage.Conditions{DoesNotExist: true})
}
writer := obj.NewWriter(ctx)
start := time.Now()
_, err = io.Copy(writer, reader)
_, err := io.Copy(writer, reader)
// Always close the writer, even if upload failed.
closeErr := writer.Close()
if err == nil {
@@ -292,28 +332,40 @@ func (h *Handler) uploadFile(ctx context.Context, path string, reader io.Reader)
// result into a writer. Returns trace.NotFound error if the recording is not
// found.
func (h *Handler) Download(ctx context.Context, sessionID session.ID, writer events.RandomAccessWriter) error {
return h.downloadFile(ctx, h.recordingPath(sessionID), writer)
return trace.Wrap(h.downloadFile(ctx, h.recordingPath(sessionID), writer))
}
// DownloadSummary downloads a session summary from a GCS bucket and writes the
// result into a writer. Returns trace.NotFound error if the recording is not
// found.
func (h *Handler) DownloadSummary(ctx context.Context, sessionID session.ID, writer events.RandomAccessWriter) error {
return h.downloadFile(ctx, h.summaryPath(sessionID), writer)
// Happy path: the final summary exists.
err := h.downloadFile(ctx, h.summaryPath(sessionID), writer)
if trace.IsNotFound(err) {
// Final summary doesn't exist, try the pending one.
err = h.downloadFile(ctx, h.pendingSummaryPath(sessionID), writer)
if trace.IsNotFound(err) {
// One more check for the final summary to prevent a race condition where
// the final one got created and the pending one got removed between the
// two checks above.
err = h.downloadFile(ctx, h.summaryPath(sessionID), writer)
}
}
return trace.Wrap(err)
}
// DownloadMetadata downloads a session's metadata from a GCS bucket and writes the
// result into a writer. Returns trace.NotFound error if the recording is not
// found.
func (h *Handler) DownloadMetadata(ctx context.Context, sessionID session.ID, writer events.RandomAccessWriter) error {
return h.downloadFile(ctx, h.metadataPath(sessionID), writer)
return trace.Wrap(h.downloadFile(ctx, h.metadataPath(sessionID), writer))
}
// DownloadThumbnail downloads a session's thumbnail from a GCS bucket and writes the
// result into a writer. Returns trace.NotFound error if the recording is not
// found.
func (h *Handler) DownloadThumbnail(ctx context.Context, sessionID session.ID, writer events.RandomAccessWriter) error {
return h.downloadFile(ctx, h.thumbnailPath(sessionID), writer)
return trace.Wrap(h.downloadFile(ctx, h.thumbnailPath(sessionID), writer))
}
func (h *Handler) downloadFile(ctx context.Context, path string, writer events.RandomAccessWriter) error {
@@ -343,6 +395,13 @@ func (h *Handler) recordingPath(sessionID session.ID) string {
return strings.TrimPrefix(path.Join(h.Path, string(sessionID)+".tar"), slash)
}
func (h *Handler) pendingSummaryPath(sessionID session.ID) string {
if h.Path == "" {
return path.Join(pendingPrefix, string(sessionID)+".summary.json")
}
return strings.TrimPrefix(path.Join(h.Path, pendingPrefix, string(sessionID)+".summary.json"), slash)
}
func (h *Handler) summaryPath(sessionID session.ID) string {
if h.Path == "" {
return string(sessionID) + ".summary.json"
+25 -10
View File
@@ -183,6 +183,26 @@ func (h *Handler) cleanupUpload(ctx context.Context, upload events.StreamUpload)
// batch delete objects to speed up the process
semCh := make(chan struct{}, maxParts)
errorsCh := make(chan error, maxParts)
// done indicates that all objects have been processed.
done := make(chan struct{})
var errors []error
// Start an error collection goroutine. Unless context gets canceled, this
// routine collects all the reported results and signals completion on the
// done channel.
go func() {
for range objects {
select {
case err := <-errorsCh:
if !trace.IsNotFound(err) {
errors = append(errors, err)
}
case <-ctx.Done():
return
}
}
done <- struct{}{}
}()
for i := range objects {
select {
case semCh <- struct{}{}:
@@ -199,16 +219,11 @@ func (h *Handler) cleanupUpload(ctx context.Context, upload events.StreamUpload)
}
}
var errors []error
for range objects {
select {
case err := <-errorsCh:
if !trace.IsNotFound(err) {
errors = append(errors, err)
}
case <-ctx.Done():
return trace.ConnectionProblem(ctx.Err(), "context closed")
}
// Wait until either ctx is canceled or we're done collecting results.
select {
case <-ctx.Done():
return trace.ConnectionProblem(ctx.Err(), "context closed")
case <-done:
}
return trace.NewAggregate(errors...)
}
+53
View File
@@ -0,0 +1,53 @@
// Teleport
// Copyright (C) 2025 Gravitational, Inc.
//
// This program is free software: you can redistribute it and/or modify
// it under the terms of the GNU Affero General Public License as published by
// the Free Software Foundation, either version 3 of the License, or
// (at your option) any later version.
//
// This program is distributed in the hope that it will be useful,
// but WITHOUT ANY WARRANTY; without even the implied warranty of
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
// GNU Affero General Public License for more details.
//
// You should have received a copy of the GNU Affero General Public License
// along with this program. If not, see <http://www.gnu.org/licenses/>.
package events
import (
"sync"
"github.com/aws/aws-sdk-go-v2/feature/s3/manager"
)
// MemBuffer is a in-memory byte buffer that implements both io.Writer and
// io.WriterAt interfaces.
type MemBuffer struct {
buf manager.WriteAtBuffer
// mu is a mutex to protect concurrent writes to the buffer. Even though the
// underlying buffer is thread-safe, we need to prevent a race condition in
// [MemBuffer.Write] between checking the length of the buffer and writing to
// it.
mu sync.Mutex
}
func (b *MemBuffer) Write(p []byte) (int, error) {
b.mu.Lock()
defer b.mu.Unlock()
return b.buf.WriteAt(p, int64(len(b.buf.Bytes())))
}
func (b *MemBuffer) WriteAt(p []byte, pos int64) (int, error) {
b.mu.Lock()
defer b.mu.Unlock()
return b.buf.WriteAt(p, pos)
}
// Bytes return the underlying byte slice.
func (b *MemBuffer) Bytes() []byte {
b.mu.Lock()
defer b.mu.Unlock()
return b.buf.Bytes()
}
+44
View File
@@ -0,0 +1,44 @@
// Teleport
// Copyright (C) 2025 Gravitational, Inc.
//
// This program is free software: you can redistribute it and/or modify
// it under the terms of the GNU Affero General Public License as published by
// the Free Software Foundation, either version 3 of the License, or
// (at your option) any later version.
//
// This program is distributed in the hope that it will be useful,
// but WITHOUT ANY WARRANTY; without even the implied warranty of
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
// GNU Affero General Public License for more details.
//
// You should have received a copy of the GNU Affero General Public License
// along with this program. If not, see <http://www.gnu.org/licenses/>.
package events
import (
"testing"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
func TestMemBuffer(t *testing.T) {
t.Parallel()
buf := &MemBuffer{}
n, err := buf.WriteAt([]byte(" likes "), 5)
require.NoError(t, err)
assert.Equal(t, 7, n)
n, err = buf.WriteAt([]byte("Alice"), 0)
require.NoError(t, err)
assert.Equal(t, 5, n)
n, err = buf.Write([]byte("Bob"))
require.NoError(t, err)
assert.Equal(t, 3, n)
assert.Equal(t, "Alice likes Bob", string(buf.Bytes()))
}
+1
View File
@@ -38,4 +38,5 @@ type s3Client interface {
ListParts(ctx context.Context, params *s3.ListPartsInput, optFns ...func(*s3.Options)) (*s3.ListPartsOutput, error)
ListMultipartUploads(ctx context.Context, params *s3.ListMultipartUploadsInput, optFns ...func(*s3.Options)) (*s3.ListMultipartUploadsOutput, error)
UploadPart(ctx context.Context, params *s3.UploadPartInput, optFns ...func(*s3.Options)) (*s3.UploadPartOutput, error)
HeadObject(ctx context.Context, params *s3.HeadObjectInput, optFns ...func(*s3.Options)) (*s3.HeadObjectOutput, error)
}
+105 -18
View File
@@ -322,36 +322,113 @@ func (h *Handler) Close() error {
// Upload reads the content of a session recording from a reader and uploads it
// to an S3 bucket. If successful, it returns URL of the uploaded object.
func (h *Handler) Upload(ctx context.Context, sessionID session.ID, reader io.Reader) (string, error) {
return h.uploadFile(ctx, h.recordingPath(sessionID), reader)
path, err := h.uploadFile(ctx, h.recordingPath(sessionID), reader)
return path, trace.Wrap(err)
}
// UploadSummary reads the content of a session summary from a reader and
// uploads it to an S3 bucket. If successful, it returns URL of the uploaded
// object.
// UploadPendingSummary reads the content of a pending session summary from a
// reader and uploads it to an S3 bucket. This function can be called multiple
// times for a given sessionID to update the state.
func (h *Handler) UploadPendingSummary(ctx context.Context, sessionID session.ID, reader io.Reader) (string, error) {
path, err := h.uploadFile(ctx, h.summaryPath(sessionID), reader, withOverwrite(), canBeOverwrittenLater())
return path, trace.Wrap(err)
}
// UploadSummary reads the content of a final version of session summary from a
// reader and uploads it to an S3 bucket. If successful, it returns URL of the
// uploaded object. This function can be called only once for a given
// sessionID; subsequent calls will return an error.
func (h *Handler) UploadSummary(ctx context.Context, sessionID session.ID, reader io.Reader) (string, error) {
return h.uploadFile(ctx, h.summaryPath(sessionID), reader)
path, err := h.uploadFile(ctx, h.summaryPath(sessionID), reader, withOverwrite())
return path, trace.Wrap(err)
}
// UploadMetadata reads the content of a session's metadata from a reader and
// uploads it to an S3 bucket. If successful, it returns URL of the uploaded
// object.
func (h *Handler) UploadMetadata(ctx context.Context, sessionID session.ID, reader io.Reader) (string, error) {
return h.uploadFile(ctx, h.metadataPath(sessionID), reader)
path, err := h.uploadFile(ctx, h.metadataPath(sessionID), reader)
return path, trace.Wrap(err)
}
// UploadThumbnail reads the content of a session's thumbnail from a reader and
// uploads it to an S3 bucket. If successful, it returns URL of the uploaded
// object.
func (h *Handler) UploadThumbnail(ctx context.Context, sessionID session.ID, reader io.Reader) (string, error) {
return h.uploadFile(ctx, h.thumbnailPath(sessionID), reader)
path, err := h.uploadFile(ctx, h.thumbnailPath(sessionID), reader)
return path, trace.Wrap(err)
}
func (h *Handler) uploadFile(ctx context.Context, path string, reader io.Reader) (string, error) {
type fileUploadConfig struct {
overwrite bool
canBeOverwrittenLater bool
}
type fileUploadOption func(*fileUploadConfig)
func withOverwrite() fileUploadOption {
return func(cfg *fileUploadConfig) {
cfg.overwrite = true
}
}
func canBeOverwrittenLater() fileUploadOption {
return func(cfg *fileUploadConfig) {
cfg.canBeOverwrittenLater = true
}
}
// overwritableKey is a metadata key that indicates the file can be later
// overwritten.
const overwritableKey = "teleport-internal-overwritable"
// uploadFile uploads a file to S3. Normally, it doesn't allow overwriting;
// this can be changed if [withOverwrite] was specified in the options. If
// file is intended to be later overwritten, [canBeOverwrittenLater] must be
// specified in the options. The file that can be overwritten is marked with
// the [overwritableKey] metadata header.
func (h *Handler) uploadFile(ctx context.Context, path string, reader io.Reader, opts ...fileUploadOption) (string, error) {
cfg := fileUploadConfig{}
for _, opt := range opts {
opt(&cfg)
}
uploadInput := &s3.PutObjectInput{
Bucket: aws.String(h.Bucket),
Key: aws.String(path),
Body: reader,
}
if cfg.canBeOverwrittenLater {
uploadInput.Metadata = map[string]string{overwritableKey: "true"}
}
if cfg.overwrite {
// File can overwrite the existing one. Let's see if the existing file's
// metadata allows it.
head, err := h.client.HeadObject(ctx, &s3.HeadObjectInput{
Bucket: aws.String(h.Bucket),
Key: aws.String(path),
})
err = awsutils.ConvertS3Error(err)
if err != nil && !trace.IsNotFound(err) {
return "", err
}
if err == nil {
// An existing file was found.
if _, overwritable := head.Metadata[overwritableKey]; !overwritable {
return "", trace.AlreadyExists("Object %q already exists and cannot be overwritten", path)
}
// Since we confirmed that this version can be overwritten, let's make
// sure no other version appeared in the meantime.
uploadInput.IfMatch = head.ETag
}
} else {
// The file shouldn't be overwritten, so we simply assert it's not there at
// all.
uploadInput.IfNoneMatch = aws.String("*")
}
if !h.Config.DisableServerSideEncryption {
uploadInput.ServerSideEncryption = awstypes.ServerSideEncryptionAwsKms
if h.Config.SSEKMSKey != "" {
@@ -372,31 +449,31 @@ func (h *Handler) uploadFile(ctx context.Context, path string, reader io.Reader)
// result into a writer. Returns trace.NotFound error if the recording is not
// found.
func (h *Handler) Download(ctx context.Context, sessionID session.ID, writer events.RandomAccessWriter) error {
return h.downloadFile(ctx, h.recordingPath(sessionID), writer)
return trace.Wrap(h.downloadOriginalFile(ctx, h.recordingPath(sessionID), writer))
}
// DownloadSummary downloads a session summary from an S3 bucket and writes the
// results into a writer. Returns trace.NotFound error if the summary is not
// found.
// DownloadSummary downloads a final session summary from an S3 bucket and
// writes the results into a writer. Returns trace.NotFound error if the
// summary is not found or is not final.
func (h *Handler) DownloadSummary(ctx context.Context, sessionID session.ID, writer events.RandomAccessWriter) error {
return h.downloadFile(ctx, h.summaryPath(sessionID), writer)
return trace.Wrap(h.downloadFile(ctx, h.summaryPath(sessionID), writer, nil /* versionID */))
}
// DownloadMetadata downloads a session's metadata from an S3 bucket and writes the
// results into a writer. Returns trace.NotFound error if the metadata is not
// found.
func (h *Handler) DownloadMetadata(ctx context.Context, sessionID session.ID, writer events.RandomAccessWriter) error {
return h.downloadFile(ctx, h.metadataPath(sessionID), writer)
return trace.Wrap(h.downloadOriginalFile(ctx, h.metadataPath(sessionID), writer))
}
// DownloadThumbnail downloads a session's thumbnail from an S3 bucket and writes the
// results into a writer. Returns trace.NotFound error if the thumbnail is not
// found.
func (h *Handler) DownloadThumbnail(ctx context.Context, sessionID session.ID, writer events.RandomAccessWriter) error {
return h.downloadFile(ctx, h.thumbnailPath(sessionID), writer)
return trace.Wrap(h.downloadOriginalFile(ctx, h.thumbnailPath(sessionID), writer))
}
func (h *Handler) downloadFile(ctx context.Context, path string, writer events.RandomAccessWriter) error {
func (h *Handler) downloadOriginalFile(ctx context.Context, path string, writer events.RandomAccessWriter) error {
// Get the oldest version of this object. This has to be done because S3
// allows overwriting objects in a bucket. To prevent corruption of recording
// data, get all versions and always return the first.
@@ -407,10 +484,20 @@ func (h *Handler) downloadFile(ctx context.Context, path string, writer events.R
h.logger.DebugContext(ctx, "Downloading file from S3", "bucket", h.Bucket, "path", path, "version_id", versionID)
_, err = h.downloader.Download(ctx, writer, &s3.GetObjectInput{
err = h.downloadFile(ctx, path, writer, aws.String(versionID))
if err != nil {
return awsutils.ConvertS3Error(err)
}
return nil
}
func (h *Handler) downloadFile(
ctx context.Context, path string, writer events.RandomAccessWriter, versionID *string,
) error {
_, err := h.downloader.Download(ctx, writer, &s3.GetObjectInput{
Bucket: aws.String(h.Bucket),
Key: aws.String(path),
VersionId: aws.String(versionID),
VersionId: versionID,
})
if err != nil {
return awsutils.ConvertS3Error(err)
+7 -1
View File
@@ -24,6 +24,7 @@ import (
"context"
"fmt"
"net/url"
"os"
"testing"
"github.com/aws/aws-sdk-go-v2/service/s3"
@@ -36,10 +37,15 @@ import (
// TestStreams tests various streaming upload scenarios
func TestStreams(t *testing.T) {
bucket := os.Getenv("TELEPORT_TEST_AUDIT_SESSIONS_S3_BUCKET")
if bucket == "" {
bucket = "teleport-unit-tests"
}
handler, err := NewHandler(context.Background(), Config{
Region: "us-west-1",
Path: "/test/",
Bucket: "teleport-unit-tests",
Bucket: bucket,
})
require.NoError(t, err)
+46 -12
View File
@@ -24,6 +24,7 @@ import (
"io"
"os"
"slices"
"strings"
"sync/atomic"
"testing"
"time"
@@ -47,16 +48,22 @@ import (
// UploadDownload tests uploads and downloads
func UploadDownload(t *testing.T, handler events.MultipartHandler) {
val := "hello, how is it going? this is the uploaded file"
ctx := t.Context()
id := session.NewID()
_, err := handler.Upload(context.TODO(), id, bytes.NewBuffer([]byte(val)))
_, err := handler.Upload(ctx, id, strings.NewReader(val))
require.NoError(t, err)
// Attempt to overwrite an existing file. This should fail.
_, err = handler.Upload(ctx, id, strings.NewReader("impostor"))
require.Error(t, err)
f, err := os.CreateTemp("", string(id))
require.NoError(t, err)
defer os.Remove(f.Name())
defer f.Close()
err = handler.Download(context.TODO(), id, f)
err = handler.Download(ctx, id, f)
require.NoError(t, err)
_, err = f.Seek(0, 0)
@@ -64,30 +71,57 @@ func UploadDownload(t *testing.T, handler events.MultipartHandler) {
data, err := io.ReadAll(f)
require.NoError(t, err)
require.Equal(t, string(data), val)
require.Equal(t, val, string(data))
}
// UploadDownloadSummary tests summary uploads and downloads
func UploadDownloadSummary(t *testing.T, handler events.MultipartHandler) {
val := "this is the summary file"
ctx := t.Context()
id := session.NewID()
_, err := handler.UploadSummary(t.Context(), id, bytes.NewBuffer([]byte(val)))
_, err := handler.UploadPendingSummary(ctx, id, strings.NewReader("pending summary"))
require.NoError(t, err)
f, err := os.CreateTemp("", string(id))
var pendingBuf events.MemBuffer
err = handler.DownloadSummary(ctx, id, &pendingBuf)
require.NoError(t, err)
defer os.Remove(f.Name())
defer f.Close()
assert.Equal(t, "pending summary", string(pendingBuf.Bytes()))
err = handler.DownloadSummary(context.TODO(), id, f)
// Override previous pending state.
_, err = handler.UploadPendingSummary(ctx, id, strings.NewReader("updated pending summary"))
require.NoError(t, err)
_, err = f.Seek(0, 0)
// Download the pending version.
var pendingBuf2 events.MemBuffer
err = handler.DownloadSummary(ctx, id, &pendingBuf2)
require.NoError(t, err)
assert.Equal(t, "updated pending summary", string(pendingBuf2.Bytes()))
// Upload the final version.
_, err = handler.UploadSummary(ctx, id, strings.NewReader("final summary"))
require.NoError(t, err)
data, err := io.ReadAll(f)
// Attempt to overwrite an existing file. This should fail.
_, err = handler.UploadSummary(ctx, id, strings.NewReader("impostor"))
require.Error(t, err)
// Download the final version.
var finalBuf events.MemBuffer
err = handler.DownloadSummary(ctx, id, &finalBuf)
require.NoError(t, err)
require.Equal(t, string(data), val)
assert.Equal(t, "final summary", string(finalBuf.Bytes()))
// Upload one more file, this time right to the final state (test if it's
// possible to upload one without a pending state).
id2 := session.NewID()
_, err = handler.UploadSummary(ctx, id2, strings.NewReader("final summary 2"))
require.NoError(t, err)
// Download the final version of the second file.
var finalBuf2 events.MemBuffer
err = handler.DownloadSummary(ctx, id2, &finalBuf2)
require.NoError(t, err)
assert.Equal(t, "final summary 2", string(finalBuf2.Bytes()))
}
// UploadDownloadMetadata tests metadata uploads and downloads
+11 -2
View File
@@ -33,10 +33,19 @@ type UploadHandler interface {
Upload(ctx context.Context, sessionID session.ID, readCloser io.Reader) (string, error)
// Download downloads a session recording and writes it to a writer.
Download(ctx context.Context, sessionID session.ID, writer RandomAccessWriter) error
// UploadSummary uploads a session summary and returns a URL with uploaded
// file in case of success.
// UploadPendingSummary uploads a pending session summary and returns a URL
// with uploaded file in case of success. This function can be called
// multiple times for a given sessionID to update the state. A pending
// session summary is any summary state that can still be later overwritten.
// It should still be contained in the same structure as the final one, but
// missing some data (in particular, the summary content itself).
UploadPendingSummary(ctx context.Context, sesisonID session.ID, readCloser io.Reader) (string, error)
// UploadSummary uploads a final session summary and returns a URL with
// uploaded file in case of success. This function can be called only once
// for a given sessionID; subsequent calls will return an error.
UploadSummary(ctx context.Context, sessionID session.ID, readCloser io.Reader) (string, error)
// DownloadSummary downloads a session summary and writes it to a writer.
// Returns a "not found" error if there's no such summary.
DownloadSummary(ctx context.Context, sessionID session.ID, writer RandomAccessWriter) error
// UploadMetadata uploads session metadata and returns a URL with the uploaded
// file in case of success. Session metadata is a file with a [recordingmetadatav1.SessionRecordingMetadata]
@@ -379,6 +379,13 @@ func (c *ErrorCountingSessionHandler) Upload(ctx context.Context, sessionID sess
return res, err
}
// UploadPendingSummary calls [c.wrapped.UploadPendingSummary] and counts the error or success.
func (c *ErrorCountingSessionHandler) UploadPendingSummary(ctx context.Context, sessionID session.ID, reader io.Reader) (string, error) {
res, err := c.wrapped.UploadPendingSummary(ctx, sessionID, reader)
c.uploads.observe(err)
return res, err
}
// UploadSummary calls [c.wrapped.UploadSummary] and counts the error or success.
func (c *ErrorCountingSessionHandler) UploadSummary(ctx context.Context, sessionID session.ID, reader io.Reader) (string, error) {
res, err := c.wrapped.UploadSummary(ctx, sessionID, reader)