Skip to content

Commit 1b4548c

Browse files
committed
support kz1 region
commit_hash:107d15f748ba862fca6f406dd3c246db78cf8351
1 parent 96de0ad commit 1b4548c

104 files changed

Lines changed: 637 additions & 122 deletions

File tree

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

.mapping.json

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -15,9 +15,14 @@
1515
"client/client.go":"security/cloudquery/cq-source-yc/client/client.go",
1616
"client/grpc.go":"security/cloudquery/cq-source-yc/client/grpc.go",
1717
"client/multiplexer.go":"security/cloudquery/cq-source-yc/client/multiplexer.go",
18+
"client/multiplexer_test.go":"security/cloudquery/cq-source-yc/client/multiplexer_test.go",
19+
"client/region.go":"security/cloudquery/cq-source-yc/client/region.go",
20+
"client/region_test.go":"security/cloudquery/cq-source-yc/client/region_test.go",
1821
"client/resolvers.go":"security/cloudquery/cq-source-yc/client/resolvers.go",
1922
"client/resolvers_test.go":"security/cloudquery/cq-source-yc/client/resolvers_test.go",
2023
"client/resourcetype.go":"security/cloudquery/cq-source-yc/client/resourcetype.go",
24+
"client/service.go":"security/cloudquery/cq-source-yc/client/service.go",
25+
"client/service_test.go":"security/cloudquery/cq-source-yc/client/service_test.go",
2126
"client/spec.go":"security/cloudquery/cq-source-yc/client/spec.go",
2227
"client/testdata/yc_compute_instances.json":"security/cloudquery/cq-source-yc/client/testdata/yc_compute_instances.json",
2328
"client/testdata/yc_kubernetes_clusters.json":"security/cloudquery/cq-source-yc/client/testdata/yc_kubernetes_clusters.json",
@@ -182,6 +187,7 @@
182187
"plugin/client.go":"security/cloudquery/cq-source-yc/plugin/client.go",
183188
"plugin/plugin.go":"security/cloudquery/cq-source-yc/plugin/plugin.go",
184189
"plugin/tables.go":"security/cloudquery/cq-source-yc/plugin/tables.go",
190+
"plugin/tables_test.go":"security/cloudquery/cq-source-yc/plugin/tables_test.go",
185191
"resources/access/access.go":"security/cloudquery/cq-source-yc/resources/access/access.go",
186192
"resources/access/access_policy.go":"security/cloudquery/cq-source-yc/resources/access/access_policy.go",
187193
"resources/access/access_policy_organizationmanager_organizations.go":"security/cloudquery/cq-source-yc/resources/access/access_policy_organizationmanager_organizations.go",

client/client.go

Lines changed: 21 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -10,6 +10,7 @@ import (
1010
"github.com/rs/zerolog"
1111
"github.com/yandex-cloud/cq-source-yc/client/yc"
1212
"github.com/yandex-cloud/cq-source-yc/client/yc/datalens"
13+
"github.com/yandex-cloud/go-genproto/yandex/cloud/endpoint"
1314
ycsdk "github.com/yandex-cloud/go-sdk"
1415
ycsdkv2 "github.com/yandex-cloud/go-sdk/v2"
1516
)
@@ -21,13 +22,18 @@ const (
2122

2223
type Client struct {
2324
hierarchy *yc.ResourceHierarchy
25+
// services the installation offers, discovered once at startup. Nil when it
26+
// was never discovered, which disables skipping altogether.
27+
services serviceSet
2428

2529
OrganizationId string
2630
CloudId string
2731
FolderId string
2832
MultiplexedResourceId string
2933
MultiplexedResourceType ResourceType
3034

35+
Region Region
36+
3137
Backend state.Client
3238
Logger zerolog.Logger
3339
SDK *ycsdk.SDK
@@ -87,6 +93,8 @@ func (c *Client) WithMultiplexedResourceId(id string) *Client {
8793
}
8894

8995
func New(ctx context.Context, logger zerolog.Logger, spec *Spec) (*Client, error) {
96+
region, _ := RegionFromEndpoint(spec.Endpoint)
97+
9098
sdk, sdkv2, err := yc.Build(ctx, logger, yc.Config{
9199
Endpoint: spec.Endpoint,
92100
UserAgent: DefaultUserAgent,
@@ -96,6 +104,11 @@ func New(ctx context.Context, logger zerolog.Logger, spec *Spec) (*Client, error
96104
if err != nil {
97105
return nil, err
98106
}
107+
// Discover available services for `sdk.KnownServices`
108+
_, err = sdk.ApiEndpoint().ApiEndpoint().List(ctx, &endpoint.ListApiEndpointsRequest{})
109+
if err != nil {
110+
return nil, fmt.Errorf("failed to discover services: %w", err)
111+
}
99112

100113
// The middleware caches IAM tokens and refreshes them before expiry.
101114
iamTokens := ycsdk.NewIAMTokenMiddleware(sdk, time.Now)
@@ -118,15 +131,22 @@ func New(ctx context.Context, logger zerolog.Logger, spec *Spec) (*Client, error
118131
SDKv2: sdkv2,
119132
Datalens: dl,
120133
Logger: logger,
134+
Region: region,
121135
}
122136

123137
hierarchy, err := yc.NewResourceHierarchy(ctx, logger, sdk, spec.OrganizationIDs, spec.CloudIDs, spec.FolderIDs)
124138
if err != nil {
125139
return nil, fmt.Errorf("fetch resource hierarchy: %w", err)
126140
}
127-
128141
client.hierarchy = hierarchy
129142

143+
client.services = newServiceSet(sdk.KnownServices())
144+
145+
client.Logger.Debug().
146+
Str("region", string(region)).
147+
Interface("services", sdk.KnownServices()).
148+
Msg("services offered by this installation")
149+
130150
if len(spec.OrganizationIDs) == 0 && len(spec.CloudIDs) == 0 && len(spec.FolderIDs) == 0 {
131151
client.Logger.Warn().Msg("no organization_ids, cloud_ids, or folder_ids specified – assuming all resources nested in all orgs")
132152
}

client/multiplexer.go

Lines changed: 57 additions & 23 deletions
Original file line numberDiff line numberDiff line change
@@ -4,43 +4,77 @@ import (
44
"github.com/cloudquery/plugin-sdk/v4/schema"
55
)
66

7-
func OrganizationMultiplex(meta schema.ClientMeta) []schema.ClientMeta {
8-
client := meta.(*Client)
9-
hierarchyItems := client.hierarchy.OrganizationRows()
7+
// available reports whether a service can be reached in the installation this
8+
// client is connected to, logging the tables we skip because it cannot.
9+
func (c *Client) available(service Service) bool {
10+
if c.serviceAvailable(service) {
11+
return true
12+
}
13+
c.Logger.Info().
14+
Str("service", string(service)).
15+
Str("region", string(c.Region)).
16+
Msg("skipping tables: this installation does not offer the service")
17+
return false
18+
}
1019

11-
var l = make([]schema.ClientMeta, len(hierarchyItems))
12-
for i, item := range hierarchyItems {
13-
l[i] = client.WithOrganization(item.Organization).WithMultiplexedResourceId(item.Organization)
20+
// GlobalMultiplex is for tables whose RPC takes no folder, cloud or
21+
// organization: it keeps the single unscoped client the scheduler would use
22+
// anyway, but drops it in regions that do not have the service.
23+
func GlobalMultiplex(service Service) schema.Multiplexer {
24+
return func(meta schema.ClientMeta) []schema.ClientMeta {
25+
client := meta.(*Client)
26+
if !client.available(service) {
27+
return nil
28+
}
29+
return []schema.ClientMeta{client}
1430
}
15-
return l
1631
}
1732

18-
func CloudMultiplex(meta schema.ClientMeta) []schema.ClientMeta {
19-
client := meta.(*Client)
20-
hierarchyItems := client.hierarchy.CloudRows()
33+
func OrganizationMultiplex(service Service) schema.Multiplexer {
34+
return func(meta schema.ClientMeta) []schema.ClientMeta {
35+
client := meta.(*Client)
36+
if !client.available(service) {
37+
return nil
38+
}
39+
hierarchyItems := client.hierarchy.OrganizationRows()
2140

22-
var l = make([]schema.ClientMeta, len(hierarchyItems))
23-
for i, item := range hierarchyItems {
24-
l[i] = client.WithOrganization(item.Organization).WithCloud(item.Cloud).WithMultiplexedResourceId(item.Cloud)
41+
var l = make([]schema.ClientMeta, len(hierarchyItems))
42+
for i, item := range hierarchyItems {
43+
l[i] = client.WithOrganization(item.Organization).WithMultiplexedResourceId(item.Organization)
44+
}
45+
return l
2546
}
26-
return l
2747
}
2848

29-
func FolderMultiplex(meta schema.ClientMeta) []schema.ClientMeta {
30-
client := meta.(*Client)
31-
hierarchyItems := client.hierarchy.FolderRows()
49+
func CloudMultiplex(service Service) schema.Multiplexer {
50+
return func(meta schema.ClientMeta) []schema.ClientMeta {
51+
client := meta.(*Client)
52+
if !client.available(service) {
53+
return nil
54+
}
55+
hierarchyItems := client.hierarchy.CloudRows()
3256

33-
var l = make([]schema.ClientMeta, len(hierarchyItems))
34-
for i, item := range hierarchyItems {
35-
l[i] = client.WithOrganization(item.Organization).WithCloud(item.Cloud).WithFolder(item.Folder).WithMultiplexedResourceId(item.Folder)
57+
var l = make([]schema.ClientMeta, len(hierarchyItems))
58+
for i, item := range hierarchyItems {
59+
l[i] = client.WithOrganization(item.Organization).WithCloud(item.Cloud).WithMultiplexedResourceId(item.Cloud)
60+
}
61+
return l
3662
}
37-
return l
3863
}
3964

40-
func PrependEmptyMultiplex(multiplexer schema.Multiplexer) schema.Multiplexer {
65+
func FolderMultiplex(service Service) schema.Multiplexer {
4166
return func(meta schema.ClientMeta) []schema.ClientMeta {
4267
client := meta.(*Client)
43-
return append([]schema.ClientMeta{client}, multiplexer(meta)...)
68+
if !client.available(service) {
69+
return nil
70+
}
71+
hierarchyItems := client.hierarchy.FolderRows()
72+
73+
var l = make([]schema.ClientMeta, len(hierarchyItems))
74+
for i, item := range hierarchyItems {
75+
l[i] = client.WithOrganization(item.Organization).WithCloud(item.Cloud).WithFolder(item.Folder).WithMultiplexedResourceId(item.Folder)
76+
}
77+
return l
4478
}
4579
}
4680

client/multiplexer_test.go

Lines changed: 69 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,69 @@
1+
package client
2+
3+
import (
4+
"testing"
5+
6+
"github.com/cloudquery/plugin-sdk/v4/schema"
7+
"github.com/rs/zerolog"
8+
"github.com/stretchr/testify/assert"
9+
)
10+
11+
// A client whose installation does not offer the service must be turned away
12+
// before the multiplexer reaches for the resource hierarchy, which is why this
13+
// passes a client that has none.
14+
func TestMultiplexSkipsUnavailableService(t *testing.T) {
15+
multiplexers := map[string]func(Service) schema.Multiplexer{
16+
"organization": OrganizationMultiplex,
17+
"cloud": CloudMultiplex,
18+
"folder": FolderMultiplex,
19+
"global": GlobalMultiplex,
20+
}
21+
22+
for name, multiplexer := range multiplexers {
23+
t.Run(name, func(t *testing.T) {
24+
c := &Client{
25+
Region: RegionKZ,
26+
Logger: zerolog.Nop(),
27+
services: newServiceSet([]string{"compute"}),
28+
}
29+
assert.Empty(t, multiplexer(ServiceBaremetal)(c))
30+
})
31+
}
32+
}
33+
34+
func TestGlobalMultiplexKeepsAvailableService(t *testing.T) {
35+
c := &Client{
36+
Region: RegionKZ,
37+
Logger: zerolog.Nop(),
38+
services: newServiceSet([]string{"compute"}),
39+
}
40+
assert.Equal(t, []schema.ClientMeta{c}, GlobalMultiplex(ServiceCompute)(c))
41+
}
42+
43+
func TestServiceAvailable(t *testing.T) {
44+
discovered := &Client{
45+
Region: RegionKZ,
46+
services: newServiceSet([]string{"compute", "vpc"}),
47+
}
48+
assert.True(t, discovered.serviceAvailable(ServiceCompute))
49+
assert.False(t, discovered.serviceAvailable(ServiceBaremetal))
50+
assert.False(t, discovered.serviceAvailable("no-such-service"))
51+
52+
t.Run("undiscovered_client_filters_nothing", func(t *testing.T) {
53+
c := &Client{Region: RegionKZ}
54+
assert.True(t, c.serviceAvailable(ServiceBaremetal))
55+
assert.True(t, c.serviceAvailable("no-such-service"))
56+
})
57+
58+
t.Run("datalens_follows_the_region", func(t *testing.T) {
59+
// DataLens is an HTTP API, so discovery never lists it and the region
60+
// has to answer for it
61+
ru := &Client{Region: RegionRU, services: newServiceSet([]string{"compute"})}
62+
kz := &Client{Region: RegionKZ, services: newServiceSet([]string{"compute"})}
63+
unknown := &Client{services: newServiceSet([]string{"compute"})}
64+
65+
assert.True(t, ru.serviceAvailable(ServiceDataLens))
66+
assert.False(t, kz.serviceAvailable(ServiceDataLens))
67+
assert.True(t, unknown.serviceAvailable(ServiceDataLens))
68+
})
69+
}

client/region.go

Lines changed: 44 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,44 @@
1+
package client
2+
3+
import (
4+
"strings"
5+
)
6+
7+
type Region string
8+
9+
const (
10+
RegionRU Region = "ru-central1"
11+
RegionKZ Region = "kz1"
12+
)
13+
14+
var RegionEndpoints = map[Region]string{
15+
RegionRU: "api.cloud.yandex.net:443",
16+
RegionKZ: "api.yandexcloud.kz:443",
17+
}
18+
19+
var regionByHost = func() map[string]Region {
20+
result := make(map[string]Region, len(RegionEndpoints))
21+
for region, endpoint := range RegionEndpoints {
22+
result[endpointHost(endpoint)] = region
23+
}
24+
return result
25+
}()
26+
27+
func RegionFromEndpoint(endpoint string) (Region, bool) {
28+
region, ok := regionByHost[endpointHost(endpoint)]
29+
return region, ok
30+
}
31+
32+
// endpointHost strips the scheme, the port and the path off an endpoint,
33+
// leaving the bare lowercase host.
34+
func endpointHost(endpoint string) string {
35+
host := endpoint
36+
if _, after, ok := strings.Cut(host, "://"); ok {
37+
host = after
38+
}
39+
host, _, _ = strings.Cut(host, "/")
40+
if i := strings.LastIndex(host, ":"); i >= 0 {
41+
host = host[:i]
42+
}
43+
return strings.ToLower(host)
44+
}

client/region_test.go

Lines changed: 43 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,43 @@
1+
package client_test
2+
3+
import (
4+
"testing"
5+
6+
"github.com/stretchr/testify/assert"
7+
"github.com/yandex-cloud/cq-source-yc/client"
8+
)
9+
10+
func TestRegionFromEndpoint(t *testing.T) {
11+
tests := []struct {
12+
endpoint string
13+
want client.Region
14+
wantOk bool
15+
}{
16+
{endpoint: "api.cloud.yandex.net:443", want: client.RegionRU, wantOk: true},
17+
{endpoint: "api.yandexcloud.kz:443", want: client.RegionKZ, wantOk: true},
18+
{endpoint: "api.cloud.yandex.net", want: client.RegionRU, wantOk: true},
19+
{endpoint: "API.YandexCloud.KZ:443", want: client.RegionKZ, wantOk: true},
20+
{endpoint: "https://api.yandexcloud.kz:443/", want: client.RegionKZ, wantOk: true},
21+
{endpoint: "api.cloud-preprod.yandex.net:443"},
22+
{endpoint: "localhost:7777"},
23+
{endpoint: ""},
24+
}
25+
26+
for _, tt := range tests {
27+
t.Run(tt.endpoint, func(t *testing.T) {
28+
region, ok := client.RegionFromEndpoint(tt.endpoint)
29+
assert.Equal(t, tt.wantOk, ok)
30+
assert.Equal(t, tt.want, region)
31+
})
32+
}
33+
}
34+
35+
// Every region we claim to know must be reachable from its own endpoint,
36+
// otherwise the availability data can never be applied to it.
37+
func TestRegionEndpointsRoundTrip(t *testing.T) {
38+
for region, endpoint := range client.RegionEndpoints {
39+
got, ok := client.RegionFromEndpoint(endpoint)
40+
assert.True(t, ok, "endpoint %q of region %q is not recognised", endpoint, region)
41+
assert.Equal(t, region, got)
42+
}
43+
}

0 commit comments

Comments
 (0)