From dff6ebdb84f5a019cd39951eac7e4d2e25c1b3eb Mon Sep 17 00:00:00 2001 From: smallchungus Date: Sat, 18 Apr 2026 09:05:02 -0400 Subject: [PATCH 1/3] chore(gmail): add google api gmail/v1 and oauth2 deps --- go.mod | 20 +++++++++++--------- go.sum | 50 ++++++++++++++++++++++++++++++-------------------- 2 files changed, 41 insertions(+), 29 deletions(-) diff --git a/go.mod b/go.mod index ce126c9..15f7bcd 100644 --- a/go.mod +++ b/go.mod @@ -4,6 +4,8 @@ go 1.25.0 require ( github.com/go-chi/chi/v5 v5.2.5 + github.com/golang-migrate/migrate/v4 v4.19.1 + github.com/google/uuid v1.6.0 github.com/jackc/pgx/v5 v5.9.1 github.com/redis/go-redis/v9 v9.18.0 github.com/testcontainers/testcontainers-go/modules/postgres v0.42.0 @@ -31,8 +33,6 @@ require ( github.com/go-logr/logr v1.4.3 // indirect github.com/go-logr/stdr v1.2.2 // indirect github.com/go-ole/go-ole v1.2.6 // indirect - github.com/golang-migrate/migrate/v4 v4.19.1 // indirect - github.com/google/uuid v1.6.0 // indirect github.com/jackc/pgerrcode v0.0.0-20220416144525-469b46aa5efa // indirect github.com/jackc/pgpassfile v1.0.0 // indirect github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761 // indirect @@ -62,14 +62,16 @@ require ( github.com/tklauser/numcpus v0.11.0 // indirect github.com/yusufpapurcu/wmi v1.2.4 // indirect go.opentelemetry.io/auto/sdk v1.2.1 // indirect - go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp v0.61.0 // indirect - go.opentelemetry.io/otel v1.41.0 // indirect - go.opentelemetry.io/otel/metric v1.41.0 // indirect - go.opentelemetry.io/otel/trace v1.41.0 // indirect + go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp v0.67.0 // indirect + go.opentelemetry.io/otel v1.43.0 // indirect + go.opentelemetry.io/otel/metric v1.43.0 // indirect + go.opentelemetry.io/otel/sdk v1.43.0 // indirect + go.opentelemetry.io/otel/sdk/metric v1.43.0 // indirect + go.opentelemetry.io/otel/trace v1.43.0 // indirect go.uber.org/atomic v1.11.0 // indirect - golang.org/x/crypto v0.48.0 // indirect - golang.org/x/sync v0.19.0 // indirect + golang.org/x/crypto v0.49.0 // indirect + golang.org/x/sync v0.20.0 // indirect golang.org/x/sys v0.42.0 // indirect - golang.org/x/text v0.34.0 // indirect + golang.org/x/text v0.35.0 // indirect gopkg.in/yaml.v3 v3.0.1 // indirect ) diff --git a/go.sum b/go.sum index 70a62e9..c98055c 100644 --- a/go.sum +++ b/go.sum @@ -31,8 +31,12 @@ github.com/davecgh/go-spew v1.1.2-0.20180830191138-d8f796af33cc h1:U9qPSI2PIWSS1 github.com/davecgh/go-spew v1.1.2-0.20180830191138-d8f796af33cc/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= github.com/dgryski/go-rendezvous v0.0.0-20200823014737-9f7001d12a5f h1:lO4WD4F/rVNCu3HqELle0jiPLLBs70cWOduZpkS1E78= github.com/dgryski/go-rendezvous v0.0.0-20200823014737-9f7001d12a5f/go.mod h1:cuUVRXasLTGF7a8hSLbxyZXjz+1KgoB3wDUb6vlszIc= +github.com/dhui/dktest v0.4.6 h1:+DPKyScKSEp3VLtbMDHcUq6V5Lm5zfZZVb0Sk7Ahom4= +github.com/dhui/dktest v0.4.6/go.mod h1:JHTSYDtKkvFNFHJKqCzVzqXecyv+tKt8EzceOmQOgbU= github.com/distribution/reference v0.6.0 h1:0IXCQ5g4/QMHHkarYzh5l+u8T3t73zM5QvfrDyIgxBk= github.com/distribution/reference v0.6.0/go.mod h1:BbU0aIcezP1/5jX/8MP0YiH4SdvB5Y4f/wlDRiLyi3E= +github.com/docker/docker v28.3.3+incompatible h1:Dypm25kh4rmk49v1eiVbsAtpAsYURjYkaKubwuBdxEI= +github.com/docker/docker v28.3.3+incompatible/go.mod h1:eEKB0N0r5NX/I1kEveEz05bcu8tLC/8azJZsviup8Sk= github.com/docker/go-connections v0.6.0 h1:LlMG9azAe1TqfR7sO+NJttz1gy6KO7VJBh+pMmjSD94= github.com/docker/go-connections v0.6.0/go.mod h1:AahvXYshr6JgfUJGdDCs2b5EZG/vmaMAntpSFH5BFKE= github.com/docker/go-units v0.5.0 h1:69rxXcBk27SvSaaxTtLh/8llcHD8vYHT7WSdRZ/jvr4= @@ -50,6 +54,8 @@ github.com/go-logr/stdr v1.2.2 h1:hSWxHoqTgW2S2qGc0LTAI563KZ5YKYRhT3MFKZMbjag= github.com/go-logr/stdr v1.2.2/go.mod h1:mMo/vtBO5dYbehREoey6XUKy/eSumjCCveDpRre4VKE= github.com/go-ole/go-ole v1.2.6 h1:/Fpf6oFPoeFik9ty7siob0G6Ke8QvQEuVcuChpwXzpY= github.com/go-ole/go-ole v1.2.6/go.mod h1:pprOEPIfldk/42T2oK7lQ4v4JSDwmV0As9GaiUsvbm0= +github.com/gogo/protobuf v1.3.2 h1:Ov1cvc58UF3b5XjBnZv7+opcTcQFZebYjWzi34vdm4Q= +github.com/gogo/protobuf v1.3.2/go.mod h1:P1XiOD3dCwIKUDQYPy72D8LYyHL2YPYrpS2s69NZV8Q= github.com/golang-migrate/migrate/v4 v4.19.1 h1:OCyb44lFuQfYXYLx1SCxPZQGU7mcaZ7gH9yH4jSFbBA= github.com/golang-migrate/migrate/v4 v4.19.1/go.mod h1:CTcgfjxhaUtsLipnLoQRWCrjYXycRz/g5+RWDuYgPrE= github.com/google/go-cmp v0.5.6/go.mod h1:v8dTdLbMG2kIc/vJvl+f65V22dbkXbowE6jgT/gNBxE= @@ -101,10 +107,14 @@ github.com/moby/sys/userns v0.1.0 h1:tVLXkFOxVu9A64/yh59slHVv9ahO9UIev4JZusOLG/g github.com/moby/sys/userns v0.1.0/go.mod h1:IHUYgu/kao6N8YZlp9Cf444ySSvCmDlmzUcYfDHOl28= github.com/moby/term v0.5.2 h1:6qk3FJAFDs6i/q3W/pQ97SX192qKfZgGjCQqfCJkgzQ= github.com/moby/term v0.5.2/go.mod h1:d3djjFCrjnB+fl8NJux+EJzu0msscUP+f8it8hPkFLc= +github.com/morikuni/aec v1.0.0 h1:nP9CBfwrvYnBRgY6qfDQkygYDmYwOilePFkwzv4dU8A= +github.com/morikuni/aec v1.0.0/go.mod h1:BbKIizmSmc5MMPqRYbxO4ZU0S0+P200+tUnFx7PXmsc= github.com/opencontainers/go-digest v1.0.0 h1:apOUWs51W5PlhuyGyz9FCeeBIOUDA/6nW8Oi/yOhh5U= github.com/opencontainers/go-digest v1.0.0/go.mod h1:0JzlMkj0TRzQZfJkVvzbP0HBR3IKzErnv2BNG4W4MAM= github.com/opencontainers/image-spec v1.1.1 h1:y0fUlFfIZhPF1W537XOLg0/fcx6zcHCJwooC2xJA040= github.com/opencontainers/image-spec v1.1.1/go.mod h1:qpqAh3Dmcf36wStyyWU+kCeDgrGnAve2nCC8+7h8Q0M= +github.com/pkg/errors v0.9.1 h1:FEBLx1zS214owpjy7qsBeixbURkuhQAwrK5UwLGTwt4= +github.com/pkg/errors v0.9.1/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0= github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= github.com/pmezard/go-difflib v1.0.1-0.20181226105442-5d4384ee4fb2 h1:Jamvg5psRIccs7FGNTlIRMkT8wgtp5eCXdBlqhYGL6U= github.com/pmezard/go-difflib v1.0.1-0.20181226105442-5d4384ee4fb2/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= @@ -141,33 +151,33 @@ github.com/zeebo/xxh3 v1.0.2 h1:xZmwmqxHZA8AI603jOQ0tMqmBr9lPeFwGg6d+xy9DC0= github.com/zeebo/xxh3 v1.0.2/go.mod h1:5NWz9Sef7zIDm2JHfFlcQvNekmcEl9ekUZQQKCYaDcA= go.opentelemetry.io/auto/sdk v1.2.1 h1:jXsnJ4Lmnqd11kwkBV2LgLoFMZKizbCi5fNZ/ipaZ64= go.opentelemetry.io/auto/sdk v1.2.1/go.mod h1:KRTj+aOaElaLi+wW1kO/DZRXwkF4C5xPbEe3ZiIhN7Y= -go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp v0.61.0 h1:F7Jx+6hwnZ41NSFTO5q4LYDtJRXBf2PD0rNBkeB/lus= -go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp v0.61.0/go.mod h1:UHB22Z8QsdRDrnAtX4PntOl36ajSxcdUMt1sF7Y6E7Q= -go.opentelemetry.io/otel v1.41.0 h1:YlEwVsGAlCvczDILpUXpIpPSL/VPugt7zHThEMLce1c= -go.opentelemetry.io/otel v1.41.0/go.mod h1:Yt4UwgEKeT05QbLwbyHXEwhnjxNO6D8L5PQP51/46dE= -go.opentelemetry.io/otel/metric v1.41.0 h1:rFnDcs4gRzBcsO9tS8LCpgR0dxg4aaxWlJxCno7JlTQ= -go.opentelemetry.io/otel/metric v1.41.0/go.mod h1:xPvCwd9pU0VN8tPZYzDZV/BMj9CM9vs00GuBjeKhJps= -go.opentelemetry.io/otel/sdk v1.36.0 h1:b6SYIuLRs88ztox4EyrvRti80uXIFy+Sqzoh9kFULbs= -go.opentelemetry.io/otel/sdk v1.36.0/go.mod h1:+lC+mTgD+MUWfjJubi2vvXWcVxyr9rmlshZni72pXeY= -go.opentelemetry.io/otel/sdk/metric v1.36.0 h1:r0ntwwGosWGaa0CrSt8cuNuTcccMXERFwHX4dThiPis= -go.opentelemetry.io/otel/sdk/metric v1.36.0/go.mod h1:qTNOhFDfKRwX0yXOqJYegL5WRaW376QbB7P4Pb0qva4= -go.opentelemetry.io/otel/trace v1.41.0 h1:Vbk2co6bhj8L59ZJ6/xFTskY+tGAbOnCtQGVVa9TIN0= -go.opentelemetry.io/otel/trace v1.41.0/go.mod h1:U1NU4ULCoxeDKc09yCWdWe+3QoyweJcISEVa1RBzOis= +go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp v0.67.0 h1:OyrsyzuttWTSur2qN/Lm0m2a8yqyIjUVBZcxFPuXq2o= +go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp v0.67.0/go.mod h1:C2NGBr+kAB4bk3xtMXfZ94gqFDtg/GkI7e9zqGh5Beg= +go.opentelemetry.io/otel v1.43.0 h1:mYIM03dnh5zfN7HautFE4ieIig9amkNANT+xcVxAj9I= +go.opentelemetry.io/otel v1.43.0/go.mod h1:JuG+u74mvjvcm8vj8pI5XiHy1zDeoCS2LB1spIq7Ay0= +go.opentelemetry.io/otel/metric v1.43.0 h1:d7638QeInOnuwOONPp4JAOGfbCEpYb+K6DVWvdxGzgM= +go.opentelemetry.io/otel/metric v1.43.0/go.mod h1:RDnPtIxvqlgO8GRW18W6Z/4P462ldprJtfxHxyKd2PY= +go.opentelemetry.io/otel/sdk v1.43.0 h1:pi5mE86i5rTeLXqoF/hhiBtUNcrAGHLKQdhg4h4V9Dg= +go.opentelemetry.io/otel/sdk v1.43.0/go.mod h1:P+IkVU3iWukmiit/Yf9AWvpyRDlUeBaRg6Y+C58QHzg= +go.opentelemetry.io/otel/sdk/metric v1.43.0 h1:S88dyqXjJkuBNLeMcVPRFXpRw2fuwdvfCGLEo89fDkw= +go.opentelemetry.io/otel/sdk/metric v1.43.0/go.mod h1:C/RJtwSEJ5hzTiUz5pXF1kILHStzb9zFlIEe85bhj6A= +go.opentelemetry.io/otel/trace v1.43.0 h1:BkNrHpup+4k4w+ZZ86CZoHHEkohws8AY+WTX09nk+3A= +go.opentelemetry.io/otel/trace v1.43.0/go.mod h1:/QJhyVBUUswCphDVxq+8mld+AvhXZLhe+8WVFxiFff0= go.uber.org/atomic v1.11.0 h1:ZvwS0R+56ePWxUNi+Atn9dWONBPp/AUETXlHW0DxSjE= go.uber.org/atomic v1.11.0/go.mod h1:LUxbIzbOniOlMKjJjyPfpl4v+PKK2cNJn91OQbhoJI0= -golang.org/x/crypto v0.48.0 h1:/VRzVqiRSggnhY7gNRxPauEQ5Drw9haKdM0jqfcCFts= -golang.org/x/crypto v0.48.0/go.mod h1:r0kV5h3qnFPlQnBSrULhlsRfryS2pmewsg+XfMgkVos= -golang.org/x/sync v0.19.0 h1:vV+1eWNmZ5geRlYjzm2adRgW2/mcpevXNg50YZtPCE4= -golang.org/x/sync v0.19.0/go.mod h1:9KTHXmSnoGruLpwFjVSX0lNNA75CykiMECbovNTZqGI= +golang.org/x/crypto v0.49.0 h1:+Ng2ULVvLHnJ/ZFEq4KdcDd/cfjrrjjNSXNzxg0Y4U4= +golang.org/x/crypto v0.49.0/go.mod h1:ErX4dUh2UM+CFYiXZRTcMpEcN8b/1gxEuv3nODoYtCA= +golang.org/x/sync v0.20.0 h1:e0PTpb7pjO8GAtTs2dQ6jYa5BWYlMuX047Dco/pItO4= +golang.org/x/sync v0.20.0/go.mod h1:9xrNwdLfx4jkKbNva9FpL6vEN7evnE43NNNJQ2LF3+0= golang.org/x/sys v0.0.0-20190916202348-b4ddaad3f8a3/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20201204225414-ed752295db88/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20210616094352-59db8d763f22/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= golang.org/x/sys v0.42.0 h1:omrd2nAlyT5ESRdCLYdm3+fMfNFE/+Rf4bDIQImRJeo= golang.org/x/sys v0.42.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw= -golang.org/x/term v0.40.0 h1:36e4zGLqU4yhjlmxEaagx2KuYbJq3EwY8K943ZsHcvg= -golang.org/x/term v0.40.0/go.mod h1:w2P8uVp06p2iyKKuvXIm7N/y0UCRt3UfJTfZ7oOpglM= -golang.org/x/text v0.34.0 h1:oL/Qq0Kdaqxa1KbNeMKwQq0reLCCaFtqu2eNuSeNHbk= -golang.org/x/text v0.34.0/go.mod h1:homfLqTYRFyVYemLBFl5GgL/DWEiH5wcsQ5gSh1yziA= +golang.org/x/term v0.41.0 h1:QCgPso/Q3RTJx2Th4bDLqML4W6iJiaXFq2/ftQF13YU= +golang.org/x/term v0.41.0/go.mod h1:3pfBgksrReYfZ5lvYM0kSO0LIkAl4Yl2bXOkKP7Ec2A= +golang.org/x/text v0.35.0 h1:JOVx6vVDFokkpaq1AEptVzLTpDe9KGpj5tR4/X+ybL8= +golang.org/x/text v0.35.0/go.mod h1:khi/HExzZJ2pGnjenulevKNX1W67CUy0AsXcNubPGCA= golang.org/x/xerrors v0.0.0-20191204190536-9bdfabe68543/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c h1:Hei/4ADfdWqJk1ZMxUNpqntNwaWcugrBjAiHlqqRiVk= From a224ec9cae963c6edd5fb570fdf2012b7a0d5adc Mon Sep 17 00:00:00 2001 From: smallchungus Date: Sat, 18 Apr 2026 09:06:28 -0400 Subject: [PATCH 2/3] feat(gmail): add LoadToken / SaveToken / savingSource helpers --- go.sum | 2 + internal/gmail/token.go | 88 +++++++++++++++++++++++++++++ internal/gmail/token_test.go | 104 +++++++++++++++++++++++++++++++++++ 3 files changed, 194 insertions(+) create mode 100644 internal/gmail/token.go create mode 100644 internal/gmail/token_test.go diff --git a/go.sum b/go.sum index c98055c..fc9d77b 100644 --- a/go.sum +++ b/go.sum @@ -167,6 +167,8 @@ go.uber.org/atomic v1.11.0 h1:ZvwS0R+56ePWxUNi+Atn9dWONBPp/AUETXlHW0DxSjE= go.uber.org/atomic v1.11.0/go.mod h1:LUxbIzbOniOlMKjJjyPfpl4v+PKK2cNJn91OQbhoJI0= golang.org/x/crypto v0.49.0 h1:+Ng2ULVvLHnJ/ZFEq4KdcDd/cfjrrjjNSXNzxg0Y4U4= golang.org/x/crypto v0.49.0/go.mod h1:ErX4dUh2UM+CFYiXZRTcMpEcN8b/1gxEuv3nODoYtCA= +golang.org/x/oauth2 v0.30.0 h1:dnDm7JmhM45NNpd8FDDeLhK6FwqbOf4MLCM9zb1BOHI= +golang.org/x/oauth2 v0.30.0/go.mod h1:B++QgG3ZKulg6sRPGD/mqlHQs5rB3Ml9erfeDY7xKlU= golang.org/x/sync v0.20.0 h1:e0PTpb7pjO8GAtTs2dQ6jYa5BWYlMuX047Dco/pItO4= golang.org/x/sync v0.20.0/go.mod h1:9xrNwdLfx4jkKbNva9FpL6vEN7evnE43NNNJQ2LF3+0= golang.org/x/sys v0.0.0-20190916202348-b4ddaad3f8a3/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= diff --git a/internal/gmail/token.go b/internal/gmail/token.go new file mode 100644 index 0000000..73fd508 --- /dev/null +++ b/internal/gmail/token.go @@ -0,0 +1,88 @@ +package gmail + +import ( + "context" + "fmt" + "time" + + "github.com/google/uuid" + "golang.org/x/oauth2" + + "github.com/smallchungus/disttaskqueue/internal/oauth" + "github.com/smallchungus/disttaskqueue/internal/store" +) + +func LoadToken(ctx context.Context, s *store.Store, userID uuid.UUID, key []byte) (*oauth2.Token, error) { + rec, err := s.GetOAuthToken(ctx, userID, "google") + if err != nil { + return nil, fmt.Errorf("load token: %w", err) + } + return decryptToken(rec.AccessCT, rec.RefreshCT, rec.ExpiresAt, key) +} + +func SaveToken(ctx context.Context, s *store.Store, userID uuid.UUID, key []byte, tok *oauth2.Token) error { + accessCT, refreshCT, err := encryptToken(tok, key) + if err != nil { + return fmt.Errorf("save token: %w", err) + } + return s.SaveOAuthToken(ctx, store.OAuthToken{ + UserID: userID, + Provider: "google", + AccessCT: accessCT, + RefreshCT: refreshCT, + ExpiresAt: tok.Expiry, + }) +} + +func encryptToken(tok *oauth2.Token, key []byte) (access, refresh []byte, err error) { + access, err = oauth.Encrypt([]byte(tok.AccessToken), key) + if err != nil { + return nil, nil, err + } + refresh, err = oauth.Encrypt([]byte(tok.RefreshToken), key) + if err != nil { + return nil, nil, err + } + return access, refresh, nil +} + +func decryptToken(access, refresh []byte, expiry time.Time, key []byte) (*oauth2.Token, error) { + a, err := oauth.Decrypt(access, key) + if err != nil { + return nil, fmt.Errorf("decrypt access: %w", err) + } + r, err := oauth.Decrypt(refresh, key) + if err != nil { + return nil, fmt.Errorf("decrypt refresh: %w", err) + } + return &oauth2.Token{ + AccessToken: string(a), + RefreshToken: string(r), + TokenType: "Bearer", + Expiry: expiry, + }, nil +} + +type savingSource struct { + base oauth2.TokenSource + save func(*oauth2.Token) error + last *oauth2.Token +} + +func newSavingSource(base oauth2.TokenSource, save func(*oauth2.Token) error, seed *oauth2.Token) oauth2.TokenSource { + return &savingSource{base: base, save: save, last: seed} +} + +func (s *savingSource) Token() (*oauth2.Token, error) { + tok, err := s.base.Token() + if err != nil { + return nil, err + } + if s.last == nil || tok.AccessToken != s.last.AccessToken { + if err := s.save(tok); err != nil { + return nil, fmt.Errorf("save: %w", err) + } + s.last = tok + } + return tok, nil +} diff --git a/internal/gmail/token_test.go b/internal/gmail/token_test.go new file mode 100644 index 0000000..f00a004 --- /dev/null +++ b/internal/gmail/token_test.go @@ -0,0 +1,104 @@ +package gmail + +import ( + "context" + "errors" + "testing" + "time" + + "github.com/google/uuid" + "golang.org/x/oauth2" + + "github.com/smallchungus/disttaskqueue/internal/store" +) + +type fakeStore struct{} + +type sequenceSource struct { + tokens []*oauth2.Token + idx int + err error +} + +func (s *sequenceSource) Token() (*oauth2.Token, error) { + if s.err != nil { + return nil, s.err + } + if s.idx >= len(s.tokens) { + return s.tokens[len(s.tokens)-1], nil + } + tok := s.tokens[s.idx] + s.idx++ + return tok, nil +} + +func TestSavingSource_SavesWhenTokenChanges(t *testing.T) { + tok1 := &oauth2.Token{AccessToken: "v1", RefreshToken: "r1", Expiry: time.Now().Add(time.Hour)} + tok2 := &oauth2.Token{AccessToken: "v2", RefreshToken: "r1", Expiry: time.Now().Add(2 * time.Hour)} + + base := &sequenceSource{tokens: []*oauth2.Token{tok1, tok2}} + var saved []*oauth2.Token + src := newSavingSource(base, func(t *oauth2.Token) error { + saved = append(saved, t) + return nil + }, tok1) + + got, err := src.Token() + if err != nil { + t.Fatal(err) + } + if got.AccessToken != "v1" { + t.Fatalf("first token: %s", got.AccessToken) + } + if len(saved) != 0 { + t.Fatalf("save called for unchanged token: %d", len(saved)) + } + + got, err = src.Token() + if err != nil { + t.Fatal(err) + } + if got.AccessToken != "v2" { + t.Fatalf("second token: %s", got.AccessToken) + } + if len(saved) != 1 || saved[0].AccessToken != "v2" { + t.Fatalf("save not called with refreshed token: %+v", saved) + } +} + +func TestSavingSource_PropagatesBaseError(t *testing.T) { + base := &sequenceSource{err: errors.New("refresh failed")} + src := newSavingSource(base, func(*oauth2.Token) error { return nil }, nil) + if _, err := src.Token(); err == nil { + t.Fatal("expected error") + } +} + +func TestEncryptToken_DecryptToken_RoundTrip(t *testing.T) { + key := make([]byte, 32) + for i := range key { + key[i] = byte(i) + } + in := &oauth2.Token{ + AccessToken: "the-access", + RefreshToken: "the-refresh", + Expiry: time.Now().Add(time.Hour).UTC().Truncate(time.Second), + } + + accessCT, refreshCT, err := encryptToken(in, key) + if err != nil { + t.Fatalf("encrypt: %v", err) + } + out, err := decryptToken(accessCT, refreshCT, in.Expiry, key) + if err != nil { + t.Fatalf("decrypt: %v", err) + } + if out.AccessToken != in.AccessToken || out.RefreshToken != in.RefreshToken { + t.Fatalf("mismatch: %+v vs %+v", out, in) + } +} + +var _ = fakeStore{} +var _ = uuid.UUID{} +var _ = context.Background() +var _ = store.OAuthToken{} From 917e43d76f5714052d5c0b7a47151a9320d960c9 Mon Sep 17 00:00:00 2001 From: smallchungus Date: Sat, 18 Apr 2026 09:08:15 -0400 Subject: [PATCH 3/3] feat(gmail): add Client.New and LatestMessageIDs via History API --- go.mod | 12 ++ go.sum | 24 ++++ internal/gmail/client.go | 112 +++++++++++++++ internal/gmail/client_integration_test.go | 160 ++++++++++++++++++++++ 4 files changed, 308 insertions(+) create mode 100644 internal/gmail/client.go create mode 100644 internal/gmail/client_integration_test.go diff --git a/go.mod b/go.mod index 15f7bcd..18362a8 100644 --- a/go.mod +++ b/go.mod @@ -10,9 +10,14 @@ require ( github.com/redis/go-redis/v9 v9.18.0 github.com/testcontainers/testcontainers-go/modules/postgres v0.42.0 github.com/testcontainers/testcontainers-go/modules/redis v0.42.0 + golang.org/x/oauth2 v0.30.0 + google.golang.org/api v0.247.0 ) require ( + cloud.google.com/go/auth v0.16.4 // indirect + cloud.google.com/go/auth/oauth2adapt v0.2.8 // indirect + cloud.google.com/go/compute/metadata v0.8.0 // indirect dario.cat/mergo v1.0.2 // indirect github.com/Azure/go-ansiterm v0.0.0-20250102033503-faa5f7b0171c // indirect github.com/Microsoft/go-winio v0.6.2 // indirect @@ -33,6 +38,9 @@ require ( github.com/go-logr/logr v1.4.3 // indirect github.com/go-logr/stdr v1.2.2 // indirect github.com/go-ole/go-ole v1.2.6 // indirect + github.com/google/s2a-go v0.1.9 // indirect + github.com/googleapis/enterprise-certificate-proxy v0.3.6 // indirect + github.com/googleapis/gax-go/v2 v2.15.0 // indirect github.com/jackc/pgerrcode v0.0.0-20220416144525-469b46aa5efa // indirect github.com/jackc/pgpassfile v1.0.0 // indirect github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761 // indirect @@ -70,8 +78,12 @@ require ( go.opentelemetry.io/otel/trace v1.43.0 // indirect go.uber.org/atomic v1.11.0 // indirect golang.org/x/crypto v0.49.0 // indirect + golang.org/x/net v0.51.0 // indirect golang.org/x/sync v0.20.0 // indirect golang.org/x/sys v0.42.0 // indirect golang.org/x/text v0.35.0 // indirect + google.golang.org/genproto/googleapis/rpc v0.0.0-20250818200422-3122310a409c // indirect + google.golang.org/grpc v1.74.2 // indirect + google.golang.org/protobuf v1.36.7 // indirect gopkg.in/yaml.v3 v3.0.1 // indirect ) diff --git a/go.sum b/go.sum index fc9d77b..ba87ef8 100644 --- a/go.sum +++ b/go.sum @@ -1,3 +1,10 @@ +cloud.google.com/go v0.121.6 h1:waZiuajrI28iAf40cWgycWNgaXPO06dupuS+sgibK6c= +cloud.google.com/go/auth v0.16.4 h1:fXOAIQmkApVvcIn7Pc2+5J8QTMVbUGLscnSVNl11su8= +cloud.google.com/go/auth v0.16.4/go.mod h1:j10ncYwjX/g3cdX7GpEzsdM+d+ZNsXAbb6qXA7p1Y5M= +cloud.google.com/go/auth/oauth2adapt v0.2.8 h1:keo8NaayQZ6wimpNSmW5OPc283g65QNIiLpZnkHRbnc= +cloud.google.com/go/auth/oauth2adapt v0.2.8/go.mod h1:XQ9y31RkqZCcwJWNSx2Xvric3RrU88hAYYbjDWYDL+c= +cloud.google.com/go/compute/metadata v0.8.0 h1:HxMRIbao8w17ZX6wBnjhcDkW6lTFpgcaobyVfZWqRLA= +cloud.google.com/go/compute/metadata v0.8.0/go.mod h1:sYOGTp851OV9bOFJ9CH7elVvyzopvWQFNNghtDQ/Biw= dario.cat/mergo v1.0.2 h1:85+piFYR1tMbRrLcDwR18y4UKJ3aH1Tbzi24VRW1TK8= dario.cat/mergo v1.0.2/go.mod h1:E/hbnu0NxMFBjpMIE34DRGLWqDy0g5FuKDhCb31ngxA= github.com/AdaLogics/go-fuzz-headers v0.0.0-20240806141605-e8a1dd7889d6 h1:He8afgbRMd7mFxO99hRNu+6tazq8nFF9lIwo9JFroBk= @@ -61,8 +68,14 @@ github.com/golang-migrate/migrate/v4 v4.19.1/go.mod h1:CTcgfjxhaUtsLipnLoQRWCrjY github.com/google/go-cmp v0.5.6/go.mod h1:v8dTdLbMG2kIc/vJvl+f65V22dbkXbowE6jgT/gNBxE= github.com/google/go-cmp v0.7.0 h1:wk8382ETsv4JYUZwIsn6YpYiWiBsYLSJiTsyBybVuN8= github.com/google/go-cmp v0.7.0/go.mod h1:pXiqmnSA92OHEEa9HXL2W4E7lf9JzCmGVUdgjX3N/iU= +github.com/google/s2a-go v0.1.9 h1:LGD7gtMgezd8a/Xak7mEWL0PjoTQFvpRudN895yqKW0= +github.com/google/s2a-go v0.1.9/go.mod h1:YA0Ei2ZQL3acow2O62kdp9UlnvMmU7kA6Eutn0dXayM= github.com/google/uuid v1.6.0 h1:NIvaJDMOsjHA8n1jAhLSgzrAzy1Hgr+hNrb57e+94F0= github.com/google/uuid v1.6.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo= +github.com/googleapis/enterprise-certificate-proxy v0.3.6 h1:GW/XbdyBFQ8Qe+YAmFU9uHLo7OnF5tL52HFAgMmyrf4= +github.com/googleapis/enterprise-certificate-proxy v0.3.6/go.mod h1:MkHOF77EYAE7qfSuSS9PU6g4Nt4e11cnsDUowfwewLA= +github.com/googleapis/gax-go/v2 v2.15.0 h1:SyjDc1mGgZU5LncH8gimWo9lW1DtIfPibOG81vgd/bo= +github.com/googleapis/gax-go/v2 v2.15.0/go.mod h1:zVVkkxAQHa1RQpg9z2AUCMnKhi0Qld9rcmyfL1OZhoc= github.com/jackc/pgerrcode v0.0.0-20220416144525-469b46aa5efa h1:s+4MhCQ6YrzisK6hFJUX53drDT4UsSW3DEhKn0ifuHw= github.com/jackc/pgerrcode v0.0.0-20220416144525-469b46aa5efa/go.mod h1:a/s9Lp5W7n/DD0VrVoyJ00FbP2ytTPDVOivvn2bMlds= github.com/jackc/pgpassfile v1.0.0 h1:/6Hmqy13Ss2zCq62VdNG8tM1wchn8zjSGOBJ6icpsIM= @@ -167,6 +180,8 @@ go.uber.org/atomic v1.11.0 h1:ZvwS0R+56ePWxUNi+Atn9dWONBPp/AUETXlHW0DxSjE= go.uber.org/atomic v1.11.0/go.mod h1:LUxbIzbOniOlMKjJjyPfpl4v+PKK2cNJn91OQbhoJI0= golang.org/x/crypto v0.49.0 h1:+Ng2ULVvLHnJ/ZFEq4KdcDd/cfjrrjjNSXNzxg0Y4U4= golang.org/x/crypto v0.49.0/go.mod h1:ErX4dUh2UM+CFYiXZRTcMpEcN8b/1gxEuv3nODoYtCA= +golang.org/x/net v0.51.0 h1:94R/GTO7mt3/4wIKpcR5gkGmRLOuE/2hNGeWq/GBIFo= +golang.org/x/net v0.51.0/go.mod h1:aamm+2QF5ogm02fjy5Bb7CQ0WMt1/WVM7FtyaTLlA9Y= golang.org/x/oauth2 v0.30.0 h1:dnDm7JmhM45NNpd8FDDeLhK6FwqbOf4MLCM9zb1BOHI= golang.org/x/oauth2 v0.30.0/go.mod h1:B++QgG3ZKulg6sRPGD/mqlHQs5rB3Ml9erfeDY7xKlU= golang.org/x/sync v0.20.0 h1:e0PTpb7pjO8GAtTs2dQ6jYa5BWYlMuX047Dco/pItO4= @@ -181,6 +196,15 @@ golang.org/x/term v0.41.0/go.mod h1:3pfBgksrReYfZ5lvYM0kSO0LIkAl4Yl2bXOkKP7Ec2A= golang.org/x/text v0.35.0 h1:JOVx6vVDFokkpaq1AEptVzLTpDe9KGpj5tR4/X+ybL8= golang.org/x/text v0.35.0/go.mod h1:khi/HExzZJ2pGnjenulevKNX1W67CUy0AsXcNubPGCA= golang.org/x/xerrors v0.0.0-20191204190536-9bdfabe68543/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= +google.golang.org/api v0.247.0 h1:tSd/e0QrUlLsrwMKmkbQhYVa109qIintOls2Wh6bngc= +google.golang.org/api v0.247.0/go.mod h1:r1qZOPmxXffXg6xS5uhx16Fa/UFY8QU/K4bfKrnvovM= +google.golang.org/genproto v0.0.0-20250603155806-513f23925822 h1:rHWScKit0gvAPuOnu87KpaYtjK5zBMLcULh7gxkCXu4= +google.golang.org/genproto/googleapis/rpc v0.0.0-20250818200422-3122310a409c h1:qXWI/sQtv5UKboZ/zUk7h+mrf/lXORyI+n9DKDAusdg= +google.golang.org/genproto/googleapis/rpc v0.0.0-20250818200422-3122310a409c/go.mod h1:gw1tLEfykwDz2ET4a12jcXt4couGAm7IwsVaTy0Sflo= +google.golang.org/grpc v1.74.2 h1:WoosgB65DlWVC9FqI82dGsZhWFNBSLjQ84bjROOpMu4= +google.golang.org/grpc v1.74.2/go.mod h1:CtQ+BGjaAIXHs/5YS3i473GqwBBa1zGQNevxdeBEXrM= +google.golang.org/protobuf v1.36.7 h1:IgrO7UwFQGJdRNXH/sQux4R1Dj1WAKcLElzeeRaXV2A= +google.golang.org/protobuf v1.36.7/go.mod h1:jduwjTPXsFjZGTmRluh+L6NjiWu7pchiJ2/5YcXBHnY= gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c h1:Hei/4ADfdWqJk1ZMxUNpqntNwaWcugrBjAiHlqqRiVk= gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c/go.mod h1:JHkPIbrfpd72SG/EVd6muEfDQjcINNoR0C8j2r3qZ4Q= diff --git a/internal/gmail/client.go b/internal/gmail/client.go new file mode 100644 index 0000000..3a50b96 --- /dev/null +++ b/internal/gmail/client.go @@ -0,0 +1,112 @@ +package gmail + +import ( + "context" + "encoding/base64" + "fmt" + "strconv" + + "github.com/google/uuid" + "golang.org/x/oauth2" + gmailapi "google.golang.org/api/gmail/v1" + "google.golang.org/api/option" + + "github.com/smallchungus/disttaskqueue/internal/store" +) + +type Config struct { + Store *store.Store + UserID uuid.UUID + EncryptionKey []byte + OAuth2 *oauth2.Config + Endpoint string +} + +type Client struct { + svc *gmailapi.Service + store *store.Store + userID uuid.UUID + key []byte +} + +func New(ctx context.Context, cfg Config) (*Client, error) { + tok, err := LoadToken(ctx, cfg.Store, cfg.UserID, cfg.EncryptionKey) + if err != nil { + return nil, err + } + + base := cfg.OAuth2.TokenSource(ctx, tok) + saving := newSavingSource(base, func(t *oauth2.Token) error { + return SaveToken(ctx, cfg.Store, cfg.UserID, cfg.EncryptionKey, t) + }, tok) + httpClient := oauth2.NewClient(ctx, saving) + + opts := []option.ClientOption{option.WithHTTPClient(httpClient)} + if cfg.Endpoint != "" { + opts = append(opts, option.WithEndpoint(cfg.Endpoint)) + } + svc, err := gmailapi.NewService(ctx, opts...) + if err != nil { + return nil, fmt.Errorf("gmail svc: %w", err) + } + return &Client{svc: svc, store: cfg.Store, userID: cfg.UserID, key: cfg.EncryptionKey}, nil +} + +func (c *Client) LatestMessageIDs(ctx context.Context, lastHistoryID string) (newIDs []string, newCursor string, err error) { + startID, err := strconv.ParseUint(lastHistoryID, 10, 64) + if err != nil { + return nil, "", fmt.Errorf("parse history id %q: %w", lastHistoryID, err) + } + + resp, err := c.svc.Users.History.List("me"). + StartHistoryId(startID). + HistoryTypes("messageAdded"). + LabelId("INBOX"). + Context(ctx). + Do() + if err != nil { + return nil, "", fmt.Errorf("history list: %w", err) + } + + for _, h := range resp.History { + for _, ma := range h.MessagesAdded { + if ma.Message == nil { + continue + } + if !hasLabel(ma.Message.LabelIds, "CATEGORY_PERSONAL") { + continue + } + newIDs = append(newIDs, ma.Message.Id) + } + } + + cursor := lastHistoryID + if resp.HistoryId != 0 { + cursor = strconv.FormatUint(resp.HistoryId, 10) + } + return newIDs, cursor, nil +} + +func (c *Client) FetchMessage(ctx context.Context, messageID string) ([]byte, error) { + msg, err := c.svc.Users.Messages.Get("me", messageID).Format("raw").Context(ctx).Do() + if err != nil { + return nil, fmt.Errorf("get message %s: %w", messageID, err) + } + raw, err := base64.URLEncoding.WithPadding(base64.NoPadding).DecodeString(msg.Raw) + if err != nil { + raw, err = base64.URLEncoding.DecodeString(msg.Raw) + if err != nil { + return nil, fmt.Errorf("decode raw: %w", err) + } + } + return raw, nil +} + +func hasLabel(labels []string, want string) bool { + for _, l := range labels { + if l == want { + return true + } + } + return false +} diff --git a/internal/gmail/client_integration_test.go b/internal/gmail/client_integration_test.go new file mode 100644 index 0000000..f129dec --- /dev/null +++ b/internal/gmail/client_integration_test.go @@ -0,0 +1,160 @@ +//go:build integration + +package gmail_test + +import ( + "context" + "encoding/base64" + "encoding/json" + "fmt" + "net/http" + "net/http/httptest" + "testing" + "time" + + "golang.org/x/oauth2" + + "github.com/smallchungus/disttaskqueue/internal/gmail" + "github.com/smallchungus/disttaskqueue/internal/store" + "github.com/smallchungus/disttaskqueue/internal/testutil" +) + +type gmailMock struct { + historyJSON string + messageJSON string + historyHits int + messageHits int +} + +func (m *gmailMock) server(t *testing.T) *httptest.Server { + t.Helper() + mux := http.NewServeMux() + mux.HandleFunc("/gmail/v1/users/me/history", func(w http.ResponseWriter, _ *http.Request) { + m.historyHits++ + w.Header().Set("Content-Type", "application/json") + _, _ = w.Write([]byte(m.historyJSON)) + }) + mux.HandleFunc("/gmail/v1/users/me/messages/", func(w http.ResponseWriter, _ *http.Request) { + m.messageHits++ + w.Header().Set("Content-Type", "application/json") + _, _ = w.Write([]byte(m.messageJSON)) + }) + srv := httptest.NewServer(mux) + t.Cleanup(srv.Close) + return srv +} + +func newKey32() []byte { + k := make([]byte, 32) + for i := range k { + k[i] = byte(i) + } + return k +} + +func setupClient(t *testing.T, m *gmailMock) (*gmail.Client, *store.Store) { + t.Helper() + pool := testutil.StartPostgres(t) + if err := store.Migrate(context.Background(), pool.Config().ConnString()); err != nil { + t.Fatalf("migrate: %v", err) + } + s := store.New(pool) + + u, err := s.CreateUser(context.Background(), fmt.Sprintf("test+%d@example.com", time.Now().UnixNano())) + if err != nil { + t.Fatal(err) + } + + key := newKey32() + tok := &oauth2.Token{ + AccessToken: "fake-access", + RefreshToken: "fake-refresh", + Expiry: time.Now().Add(1 * time.Hour), + } + if err := gmail.SaveToken(context.Background(), s, u.ID, key, tok); err != nil { + t.Fatal(err) + } + + srv := m.server(t) + cfg := &oauth2.Config{ClientID: "x", ClientSecret: "y"} + c, err := gmail.New(context.Background(), gmail.Config{ + Store: s, + UserID: u.ID, + EncryptionKey: key, + OAuth2: cfg, + Endpoint: srv.URL, + }) + if err != nil { + t.Fatalf("new client: %v", err) + } + return c, s +} + +func TestLatestMessageIDs_ReturnsInboxPrimaryAdds(t *testing.T) { + histResp := map[string]any{ + "history": []map[string]any{ + { + "id": "100", + "messagesAdded": []map[string]any{ + {"message": map[string]any{"id": "m1", "labelIds": []string{"INBOX", "CATEGORY_PERSONAL"}}}, + {"message": map[string]any{"id": "m2", "labelIds": []string{"INBOX", "CATEGORY_PROMOTIONS"}}}, + {"message": map[string]any{"id": "m3", "labelIds": []string{"INBOX", "CATEGORY_PERSONAL"}}}, + }, + }, + }, + "historyId": "150", + } + b, _ := json.Marshal(histResp) + + c, _ := setupClient(t, &gmailMock{historyJSON: string(b)}) + + ids, cursor, err := c.LatestMessageIDs(context.Background(), "1") + if err != nil { + t.Fatalf("call: %v", err) + } + if len(ids) != 2 || ids[0] != "m1" || ids[1] != "m3" { + t.Fatalf("ids: %v, want [m1 m3]", ids) + } + if cursor != "150" { + t.Fatalf("cursor: %q, want 150", cursor) + } +} + +func TestLatestMessageIDs_EmptyHistory(t *testing.T) { + histResp := map[string]any{"historyId": "200"} + b, _ := json.Marshal(histResp) + + c, _ := setupClient(t, &gmailMock{historyJSON: string(b)}) + + ids, cursor, err := c.LatestMessageIDs(context.Background(), "1") + if err != nil { + t.Fatalf("call: %v", err) + } + if len(ids) != 0 { + t.Fatalf("ids: %v, want empty", ids) + } + if cursor != "200" { + t.Fatalf("cursor: %q, want 200", cursor) + } +} + +func TestFetchMessage_DecodesRawBase64(t *testing.T) { + rawMime := "From: alice@example.com\r\nSubject: hi\r\n\r\nbody" + rawB64 := base64URL(rawMime) + msgResp := map[string]any{"raw": rawB64} + b, _ := json.Marshal(msgResp) + + c, _ := setupClient(t, &gmailMock{messageJSON: string(b)}) + + got, err := c.FetchMessage(context.Background(), "m1") + if err != nil { + t.Fatalf("fetch: %v", err) + } + if string(got) != rawMime { + t.Fatalf("got %q, want %q", got, rawMime) + } +} + +func base64URL(s string) string { + return base64.URLEncoding.WithPadding(base64.NoPadding).EncodeToString([]byte(s)) +}