From 07257d7f872f8ebf09ae83ab1aba0ebec1f8dac7 Mon Sep 17 00:00:00 2001 From: "dependabot[bot]" <49699333+dependabot[bot]@users.noreply.github.com> Date: Mon, 19 May 2025 21:05:09 +0000 Subject: [PATCH] Bump github.com/conductorone/baton-sdk from 0.3.7 to 0.3.8 Bumps [github.com/conductorone/baton-sdk](https://github.com/conductorone/baton-sdk) from 0.3.7 to 0.3.8. - [Release notes](https://github.com/conductorone/baton-sdk/releases) - [Commits](https://github.com/conductorone/baton-sdk/compare/v0.3.7...v0.3.8) --- updated-dependencies: - dependency-name: github.com/conductorone/baton-sdk dependency-version: 0.3.8 dependency-type: direct:production update-type: version-update:semver-patch ... Signed-off-by: dependabot[bot] --- go.mod | 2 +- go.sum | 4 +- .../pkg/connectorbuilder/connectorbuilder.go | 2 +- .../pkg/connectorstore/connectorstore.go | 1 + .../baton-sdk/pkg/dotc1z/sync_runs.go | 13 ++ .../baton-sdk/pkg/sdk/empty_connector.go | 137 +++++++++++++---- .../conductorone/baton-sdk/pkg/sdk/version.go | 2 +- .../conductorone/baton-sdk/pkg/sync/state.go | 34 ++-- .../conductorone/baton-sdk/pkg/sync/syncer.go | 134 ++++++++++++++-- .../baton-sdk/pkg/synccompactor/compactor.go | 145 +++++++++--------- .../baton-sdk/pkg/tasks/local/compactor.go | 6 +- vendor/modules.txt | 2 +- 12 files changed, 355 insertions(+), 127 deletions(-) diff --git a/go.mod b/go.mod index bf2b9257..2af6117e 100644 --- a/go.mod +++ b/go.mod @@ -3,7 +3,7 @@ module github.com/conductorone/baton-sql go 1.24 require ( - github.com/conductorone/baton-sdk v0.3.7 + github.com/conductorone/baton-sdk v0.3.8 github.com/elliotchance/phpserialize v1.4.0 github.com/ennyjfrick/ruleguard-logfatal v0.0.2 github.com/go-sql-driver/mysql v1.9.2 diff --git a/go.sum b/go.sum index c96e4612..bad6483b 100644 --- a/go.sum +++ b/go.sum @@ -74,8 +74,8 @@ github.com/cenkalti/backoff/v4 v4.3.0/go.mod h1:Y3VNntkOUPxTVeUxJ/G5vcM//AlwfmyY github.com/census-instrumentation/opencensus-proto v0.2.1/go.mod h1:f6KPmirojxKA12rnyqOA5BBL4O983OfeGPqjHWSTneU= github.com/client9/misspell v0.3.4/go.mod h1:qj6jICC3Q7zFZvVWo7KLAzC3yx5G7kyvSDkc90ppPyw= github.com/cncf/udpa/go v0.0.0-20191209042840-269d4d468f6f/go.mod h1:M8M6+tZqaGXZJjfX53e64911xZQV5JYwmTeXPW+k8Sc= -github.com/conductorone/baton-sdk v0.3.7 h1:A6GNQtzsqPrupGsdFxA1/1kV4dMrGpP9n1r3Ahqpsjw= -github.com/conductorone/baton-sdk v0.3.7/go.mod h1:lWZHgu025Rsgs5jvBrhilGti0zWF2+YfaFY/bWOS/g0= +github.com/conductorone/baton-sdk v0.3.8 h1:AuRBGVbgw8TcDPix7Ho7qwMucsPqP8hIh2lABUbY7Xg= +github.com/conductorone/baton-sdk v0.3.8/go.mod h1:lWZHgu025Rsgs5jvBrhilGti0zWF2+YfaFY/bWOS/g0= github.com/conductorone/dpop v0.2.4 h1:PaiDOX1gAIXtOJPxXf08GsGkpCuT/iECEjSJzLpi0zU= github.com/conductorone/dpop v0.2.4/go.mod h1:gyo8TtzB9SCFCsjsICH4IaLZ7y64CcrDXMOPBwfq/3s= github.com/conductorone/dpop/integrations/dpop_grpc v0.2.4 h1:lYxYi9/WTSL9sE96CO0QF2BY3kehs8dTTApI134TGCA= diff --git a/vendor/github.com/conductorone/baton-sdk/pkg/connectorbuilder/connectorbuilder.go b/vendor/github.com/conductorone/baton-sdk/pkg/connectorbuilder/connectorbuilder.go index f28b8908..b5f0cf77 100644 --- a/vendor/github.com/conductorone/baton-sdk/pkg/connectorbuilder/connectorbuilder.go +++ b/vendor/github.com/conductorone/baton-sdk/pkg/connectorbuilder/connectorbuilder.go @@ -650,7 +650,7 @@ func (b *builderImpl) GetResource(ctx context.Context, request *v2.ResourceGette rb, ok := b.resourceTargetedSyncers[resourceType] if !ok { b.m.RecordTaskFailure(ctx, tt, b.nowFunc().Sub(start)) - return nil, fmt.Errorf("error: get resource with unknown resource type %s", resourceType) + return nil, status.Errorf(codes.Unimplemented, "error: get resource with unknown resource type %s", resourceType) } resource, annos, err := rb.Get(ctx, request.GetResourceId(), request.GetParentResourceId()) diff --git a/vendor/github.com/conductorone/baton-sdk/pkg/connectorstore/connectorstore.go b/vendor/github.com/conductorone/baton-sdk/pkg/connectorstore/connectorstore.go index c2482e81..71db3b92 100644 --- a/vendor/github.com/conductorone/baton-sdk/pkg/connectorstore/connectorstore.go +++ b/vendor/github.com/conductorone/baton-sdk/pkg/connectorstore/connectorstore.go @@ -37,6 +37,7 @@ type Writer interface { StartSync(ctx context.Context) (string, bool, error) StartNewSync(ctx context.Context) (string, error) StartNewSyncV2(ctx context.Context, syncType string, parentSyncID string) (string, error) + SetCurrentSync(ctx context.Context, syncID string) error CurrentSyncStep(ctx context.Context) (string, error) CheckpointSync(ctx context.Context, syncToken string) error EndSync(ctx context.Context) error diff --git a/vendor/github.com/conductorone/baton-sdk/pkg/dotc1z/sync_runs.go b/vendor/github.com/conductorone/baton-sdk/pkg/dotc1z/sync_runs.go index cb32c0f9..96fbfc9a 100644 --- a/vendor/github.com/conductorone/baton-sdk/pkg/dotc1z/sync_runs.go +++ b/vendor/github.com/conductorone/baton-sdk/pkg/dotc1z/sync_runs.go @@ -340,6 +340,19 @@ func (c *C1File) getCurrentSync(ctx context.Context) (*syncRun, error) { return c.getSync(ctx, c.currentSyncID) } +func (c *C1File) SetCurrentSync(ctx context.Context, syncID string) error { + ctx, span := tracer.Start(ctx, "C1File.SetCurrentSync") + defer span.End() + + _, err := c.getSync(ctx, syncID) + if err != nil { + return err + } + + c.currentSyncID = syncID + return nil +} + func (c *C1File) CheckpointSync(ctx context.Context, syncToken string) error { ctx, span := tracer.Start(ctx, "C1File.CheckpointSync") defer span.End() diff --git a/vendor/github.com/conductorone/baton-sdk/pkg/sdk/empty_connector.go b/vendor/github.com/conductorone/baton-sdk/pkg/sdk/empty_connector.go index 19b7dd35..20cbcd3f 100644 --- a/vendor/github.com/conductorone/baton-sdk/pkg/sdk/empty_connector.go +++ b/vendor/github.com/conductorone/baton-sdk/pkg/sdk/empty_connector.go @@ -4,71 +4,156 @@ import ( "context" v2 "github.com/conductorone/baton-sdk/pb/c1/connector/v2" + "google.golang.org/grpc" + "google.golang.org/grpc/codes" + "google.golang.org/grpc/status" ) type emptyConnector struct{} // GetAsset gets an asset. -func (n *emptyConnector) GetAsset(request *v2.AssetServiceGetAssetRequest, server v2.AssetService_GetAssetServer) error { - err := server.Send(&v2.AssetServiceGetAssetResponse{ - Msg: &v2.AssetServiceGetAssetResponse_Metadata_{ - Metadata: &v2.AssetServiceGetAssetResponse_Metadata{ContentType: "application/example"}, - }, - }) - if err != nil { - return err - } - - err = server.Send(&v2.AssetServiceGetAssetResponse{ - Msg: &v2.AssetServiceGetAssetResponse_Data_{ - Data: &v2.AssetServiceGetAssetResponse_Data{Data: nil}, - }, - }) - if err != nil { - return err - } - - return nil +func (n *emptyConnector) GetAsset(_ context.Context, request *v2.AssetServiceGetAssetRequest, opts ...grpc.CallOption) (grpc.ServerStreamingClient[v2.AssetServiceGetAssetResponse], error) { + return nil, status.Errorf(codes.NotFound, "empty connector") } // ListResourceTypes returns a list of resource types. -func (n *emptyConnector) ListResourceTypes(ctx context.Context, request *v2.ResourceTypesServiceListResourceTypesRequest) (*v2.ResourceTypesServiceListResourceTypesResponse, error) { +func (n *emptyConnector) ListResourceTypes( + ctx context.Context, + request *v2.ResourceTypesServiceListResourceTypesRequest, + opts ...grpc.CallOption, +) (*v2.ResourceTypesServiceListResourceTypesResponse, error) { return &v2.ResourceTypesServiceListResourceTypesResponse{ List: []*v2.ResourceType{}, }, nil } // ListResources returns a list of resources. -func (n *emptyConnector) ListResources(ctx context.Context, request *v2.ResourcesServiceListResourcesRequest) (*v2.ResourcesServiceListResourcesResponse, error) { +func (n *emptyConnector) ListResources(ctx context.Context, request *v2.ResourcesServiceListResourcesRequest, opts ...grpc.CallOption) (*v2.ResourcesServiceListResourcesResponse, error) { return &v2.ResourcesServiceListResourcesResponse{ List: []*v2.Resource{}, }, nil } +func (n *emptyConnector) GetResource( + ctx context.Context, + request *v2.ResourceGetterServiceGetResourceRequest, + opts ...grpc.CallOption, +) (*v2.ResourceGetterServiceGetResourceResponse, error) { + return nil, status.Errorf(codes.NotFound, "empty connector") +} + // ListEntitlements returns a list of entitlements. -func (n *emptyConnector) ListEntitlements(ctx context.Context, request *v2.EntitlementsServiceListEntitlementsRequest) (*v2.EntitlementsServiceListEntitlementsResponse, error) { +func (n *emptyConnector) ListEntitlements( + ctx context.Context, + request *v2.EntitlementsServiceListEntitlementsRequest, + opts ...grpc.CallOption, +) (*v2.EntitlementsServiceListEntitlementsResponse, error) { return &v2.EntitlementsServiceListEntitlementsResponse{ List: []*v2.Entitlement{}, }, nil } // ListGrants returns a list of grants. -func (n *emptyConnector) ListGrants(ctx context.Context, request *v2.GrantsServiceListGrantsRequest) (*v2.GrantsServiceListGrantsResponse, error) { +func (n *emptyConnector) ListGrants(ctx context.Context, request *v2.GrantsServiceListGrantsRequest, opts ...grpc.CallOption) (*v2.GrantsServiceListGrantsResponse, error) { return &v2.GrantsServiceListGrantsResponse{ List: []*v2.Grant{}, }, nil } +func (n *emptyConnector) Grant(ctx context.Context, request *v2.GrantManagerServiceGrantRequest, opts ...grpc.CallOption) (*v2.GrantManagerServiceGrantResponse, error) { + return nil, status.Errorf(codes.Unimplemented, "empty connector") +} + +func (n *emptyConnector) Revoke(ctx context.Context, request *v2.GrantManagerServiceRevokeRequest, opts ...grpc.CallOption) (*v2.GrantManagerServiceRevokeResponse, error) { + return nil, status.Errorf(codes.Unimplemented, "empty connector") +} + // GetMetadata returns a connector metadata. -func (n *emptyConnector) GetMetadata(ctx context.Context, request *v2.ConnectorServiceGetMetadataRequest) (*v2.ConnectorServiceGetMetadataResponse, error) { +func (n *emptyConnector) GetMetadata(ctx context.Context, request *v2.ConnectorServiceGetMetadataRequest, opts ...grpc.CallOption) (*v2.ConnectorServiceGetMetadataResponse, error) { return &v2.ConnectorServiceGetMetadataResponse{Metadata: &v2.ConnectorMetadata{}}, nil } // Validate is called by the connector framework to validate the correct response. -func (n *emptyConnector) Validate(ctx context.Context, request *v2.ConnectorServiceValidateRequest) (*v2.ConnectorServiceValidateResponse, error) { +func (n *emptyConnector) Validate(ctx context.Context, request *v2.ConnectorServiceValidateRequest, opts ...grpc.CallOption) (*v2.ConnectorServiceValidateResponse, error) { return &v2.ConnectorServiceValidateResponse{}, nil } +func (n *emptyConnector) BulkCreateTickets(ctx context.Context, request *v2.TicketsServiceBulkCreateTicketsRequest, opts ...grpc.CallOption) (*v2.TicketsServiceBulkCreateTicketsResponse, error) { + return nil, status.Errorf(codes.Unimplemented, "empty connector") +} + +func (n *emptyConnector) BulkGetTickets(ctx context.Context, request *v2.TicketsServiceBulkGetTicketsRequest, opts ...grpc.CallOption) (*v2.TicketsServiceBulkGetTicketsResponse, error) { + return &v2.TicketsServiceBulkGetTicketsResponse{ + Tickets: []*v2.TicketsServiceGetTicketResponse{}, + }, nil +} + +func (n *emptyConnector) CreateTicket(ctx context.Context, request *v2.TicketsServiceCreateTicketRequest, opts ...grpc.CallOption) (*v2.TicketsServiceCreateTicketResponse, error) { + return nil, status.Errorf(codes.Unimplemented, "empty connector") +} + +func (n *emptyConnector) GetTicket(ctx context.Context, request *v2.TicketsServiceGetTicketRequest, opts ...grpc.CallOption) (*v2.TicketsServiceGetTicketResponse, error) { + return nil, status.Errorf(codes.NotFound, "empty connector") +} + +func (n *emptyConnector) ListTicketSchemas(ctx context.Context, request *v2.TicketsServiceListTicketSchemasRequest, opts ...grpc.CallOption) (*v2.TicketsServiceListTicketSchemasResponse, error) { + return &v2.TicketsServiceListTicketSchemasResponse{ + List: []*v2.TicketSchema{}, + }, nil +} + +func (n *emptyConnector) GetTicketSchema(ctx context.Context, request *v2.TicketsServiceGetTicketSchemaRequest, opts ...grpc.CallOption) (*v2.TicketsServiceGetTicketSchemaResponse, error) { + return nil, status.Errorf(codes.NotFound, "empty connector") +} + +func (n *emptyConnector) Cleanup(ctx context.Context, request *v2.ConnectorServiceCleanupRequest, opts ...grpc.CallOption) (*v2.ConnectorServiceCleanupResponse, error) { + return &v2.ConnectorServiceCleanupResponse{}, nil +} + +func (n *emptyConnector) CreateAccount(ctx context.Context, request *v2.CreateAccountRequest, opts ...grpc.CallOption) (*v2.CreateAccountResponse, error) { + return nil, status.Errorf(codes.Unimplemented, "empty connector") +} + +func (n *emptyConnector) RotateCredential(ctx context.Context, request *v2.RotateCredentialRequest, opts ...grpc.CallOption) (*v2.RotateCredentialResponse, error) { + return nil, status.Errorf(codes.Unimplemented, "empty connector") +} + +func (n *emptyConnector) CreateResource(ctx context.Context, request *v2.CreateResourceRequest, opts ...grpc.CallOption) (*v2.CreateResourceResponse, error) { + return nil, status.Errorf(codes.Unimplemented, "empty connector") +} + +func (n *emptyConnector) DeleteResource(ctx context.Context, request *v2.DeleteResourceRequest, opts ...grpc.CallOption) (*v2.DeleteResourceResponse, error) { + return nil, status.Errorf(codes.Unimplemented, "empty connector") +} + +func (n *emptyConnector) DeleteResourceV2(ctx context.Context, request *v2.DeleteResourceV2Request, opts ...grpc.CallOption) (*v2.DeleteResourceV2Response, error) { + return nil, status.Errorf(codes.Unimplemented, "empty connector") +} + +func (n *emptyConnector) GetActionSchema(ctx context.Context, request *v2.GetActionSchemaRequest, opts ...grpc.CallOption) (*v2.GetActionSchemaResponse, error) { + return nil, status.Errorf(codes.NotFound, "empty connector") +} + +func (n *emptyConnector) GetActionStatus(ctx context.Context, request *v2.GetActionStatusRequest, opts ...grpc.CallOption) (*v2.GetActionStatusResponse, error) { + return nil, status.Errorf(codes.NotFound, "empty connector") +} + +func (n *emptyConnector) InvokeAction(ctx context.Context, request *v2.InvokeActionRequest, opts ...grpc.CallOption) (*v2.InvokeActionResponse, error) { + return nil, status.Errorf(codes.Unimplemented, "empty connector") +} + +func (n *emptyConnector) ListActionSchemas(ctx context.Context, request *v2.ListActionSchemasRequest, opts ...grpc.CallOption) (*v2.ListActionSchemasResponse, error) { + return &v2.ListActionSchemasResponse{ + Schemas: []*v2.BatonActionSchema{}, + }, nil +} + +func (n *emptyConnector) ListEvents(ctx context.Context, request *v2.ListEventsRequest, opts ...grpc.CallOption) (*v2.ListEventsResponse, error) { + return &v2.ListEventsResponse{ + Events: []*v2.Event{}, + }, nil +} + // NewEmptyConnector returns a new emptyConnector. func NewEmptyConnector() (*emptyConnector, error) { return &emptyConnector{}, nil diff --git a/vendor/github.com/conductorone/baton-sdk/pkg/sdk/version.go b/vendor/github.com/conductorone/baton-sdk/pkg/sdk/version.go index 0fea9493..35c99ac2 100644 --- a/vendor/github.com/conductorone/baton-sdk/pkg/sdk/version.go +++ b/vendor/github.com/conductorone/baton-sdk/pkg/sdk/version.go @@ -1,3 +1,3 @@ package sdk -const Version = "v0.3.6" +const Version = "v0.3.7" diff --git a/vendor/github.com/conductorone/baton-sdk/pkg/sync/state.go b/vendor/github.com/conductorone/baton-sdk/pkg/sync/state.go index df137208..e8dc6fe6 100644 --- a/vendor/github.com/conductorone/baton-sdk/pkg/sync/state.go +++ b/vendor/github.com/conductorone/baton-sdk/pkg/sync/state.go @@ -29,6 +29,8 @@ type State interface { SetNeedsExpansion() HasExternalResourcesGrants() bool SetHasExternalResourcesGrants() + ShouldFetchRelatedResources() bool + SetShouldFetchRelatedResources() } // ActionOp represents a sync operation. @@ -129,22 +131,24 @@ type Action struct { // state is an object used for tracking the current status of a connector sync. It operates like a stack. type state struct { - mtx sync.RWMutex - actions []Action - currentAction *Action - entitlementGraph *expand.EntitlementGraph - needsExpansion bool - hasExternalResourceGrants bool + mtx sync.RWMutex + actions []Action + currentAction *Action + entitlementGraph *expand.EntitlementGraph + needsExpansion bool + hasExternalResourceGrants bool + shouldFetchRelatedResources bool } // serializedToken is used to serialize the token to JSON. This separate object is used to avoid having exported fields // on the object used externally. We should interface this, probably. type serializedToken struct { - Actions []Action `json:"actions"` - CurrentAction *Action `json:"current_action"` - NeedsExpansion bool `json:"needs_expansion"` - EntitlementGraph *expand.EntitlementGraph `json:"entitlement_graph"` - HasExternalResourceGrants bool `json:"has_external_resource_grants"` + Actions []Action `json:"actions"` + CurrentAction *Action `json:"current_action"` + NeedsExpansion bool `json:"needs_expansion"` + EntitlementGraph *expand.EntitlementGraph `json:"entitlement_graph"` + HasExternalResourceGrants bool `json:"has_external_resource_grants"` + ShouldFetchRelatedResources bool `json:"should_fetch_related_resources"` } // push adds a new action to the stack. If there is no current state, the action is directly set to current, else @@ -286,6 +290,14 @@ func (st *state) SetHasExternalResourcesGrants() { st.hasExternalResourceGrants = true } +func (st *state) ShouldFetchRelatedResources() bool { + return st.shouldFetchRelatedResources +} + +func (st *state) SetShouldFetchRelatedResources() { + st.shouldFetchRelatedResources = true +} + // PageToken returns the page token for the current action. func (st *state) PageToken(ctx context.Context) string { c := st.Current() diff --git a/vendor/github.com/conductorone/baton-sdk/pkg/sync/syncer.go b/vendor/github.com/conductorone/baton-sdk/pkg/sync/syncer.go index 7022b979..c7c82247 100644 --- a/vendor/github.com/conductorone/baton-sdk/pkg/sync/syncer.go +++ b/vendor/github.com/conductorone/baton-sdk/pkg/sync/syncer.go @@ -204,8 +204,9 @@ type syncer struct { lastCheckPointTime time.Time counts *ProgressCounts targetedSyncResourceIDs []string - - skipEGForResourceType map[string]bool + onlyExpandGrants bool + syncID string + skipEGForResourceType map[string]bool } const minCheckpointInterval = 10 * time.Second @@ -262,6 +263,14 @@ func (s *syncer) startOrResumeSync(ctx context.Context) (string, bool, error) { // If no targetedSyncResourceIDs, find the most recent sync and resume it (regardless of partial or full). // If targetedSyncResourceIDs, start a new partial sync. Use the most recent completed sync as the parent sync ID (if it exists). + if s.syncID != "" { + err := s.store.SetCurrentSync(ctx, s.syncID) + if err != nil { + return "", false, err + } + return s.syncID, false, nil + } + var syncID string var newSync bool var err error @@ -409,6 +418,7 @@ func (s *syncer) Sync(ctx context.Context) error { ParentResourceTypeID: r.GetParentResourceId().GetResourceType(), }) } + s.state.SetShouldFetchRelatedResources() s.state.PushAction(ctx, Action{Op: SyncResourceTypesOp}) err = s.Checkpoint(ctx, true) if err != nil { @@ -424,6 +434,14 @@ func (s *syncer) Sync(ctx context.Context) error { if s.externalResourceReader != nil { s.state.PushAction(ctx, Action{Op: SyncExternalResourcesOp}) } + if s.onlyExpandGrants { + s.state.SetNeedsExpansion() + err = s.Checkpoint(ctx, true) + if err != nil { + return err + } + continue + } s.state.PushAction(ctx, Action{Op: SyncGrantsOp}) s.state.PushAction(ctx, Action{Op: SyncEntitlementsOp}) s.state.PushAction(ctx, Action{Op: SyncResourcesOp}) @@ -660,6 +678,31 @@ func (s *syncer) getSubResources(ctx context.Context, parent *v2.Resource) error return nil } +func (s *syncer) getResourceFromConnector(ctx context.Context, resourceID *v2.ResourceId, parentResourceID *v2.ResourceId) (*v2.Resource, error) { + ctx, span := tracer.Start(ctx, "syncer.getResource") + defer span.End() + + resourceResp, err := s.connector.GetResource(ctx, + &v2.ResourceGetterServiceGetResourceRequest{ + ResourceId: resourceID, + ParentResourceId: parentResourceID, + }, + ) + if err == nil { + return resourceResp.Resource, nil + } + l := ctxzap.Extract(ctx) + if status.Code(err) == codes.NotFound { + l.Warn("skipping resource due to not found", zap.String("resource_id", resourceID.GetResource()), zap.String("resource_type_id", resourceID.GetResourceType())) + return nil, nil + } + if status.Code(err) == codes.Unimplemented { + l.Warn("skipping resource due to unimplemented connector", zap.String("resource_id", resourceID.GetResource()), zap.String("resource_type_id", resourceID.GetResourceType())) + return nil, nil + } + return nil, err +} + func (s *syncer) SyncTargetedResource(ctx context.Context) error { ctx, span := tracer.Start(ctx, "syncer.SyncTargetedResource") defer span.End() @@ -680,21 +723,22 @@ func (s *syncer) SyncTargetedResource(ctx context.Context) error { } } - resourceResp, err := s.connector.GetResource(ctx, - &v2.ResourceGetterServiceGetResourceRequest{ - ResourceId: &v2.ResourceId{ - ResourceType: resourceTypeID, - Resource: resourceID, - }, - ParentResourceId: prID, - }, - ) + resource, err := s.getResourceFromConnector(ctx, &v2.ResourceId{ + ResourceType: resourceTypeID, + Resource: resourceID, + }, prID) if err != nil { return err } + // If getResource encounters not found or unimplemented, it returns a nil resource and nil error. + if resource == nil { + s.state.FinishAction(ctx) + return nil + } + // Save our resource in the DB - if err := s.store.PutResources(ctx, resourceResp.Resource); err != nil { + if err := s.store.PutResources(ctx, resource); err != nil { return err } @@ -714,7 +758,7 @@ func (s *syncer) SyncTargetedResource(ctx context.Context) error { ResourceID: resourceID, }) - err = s.getSubResources(ctx, resourceResp.Resource) + err = s.getSubResources(ctx, resource) if err != nil { return err } @@ -1515,6 +1559,7 @@ func (s *syncer) syncGrantsForResource(ctx context.Context, resourceID *v2.Resou // We want to process any grants from the previous sync first so that if there is a conflict, the newer data takes precedence grants = append(grants, resp.List...) + l := ctxzap.Extract(ctx) for _, grant := range grants { grantAnnos := annotations.Annotations(grant.GetAnnotations()) if grantAnnos.Contains(&v2.GrantExpandable{}) { @@ -1523,6 +1568,57 @@ func (s *syncer) syncGrantsForResource(ctx context.Context, resourceID *v2.Resou if grantAnnos.ContainsAny(&v2.ExternalResourceMatchAll{}, &v2.ExternalResourceMatch{}, &v2.ExternalResourceMatchID{}) { s.state.SetHasExternalResourcesGrants() } + + if !s.state.ShouldFetchRelatedResources() { + continue + } + // Some connectors emit grants for other resources. If we're doing a partial sync, check if it exists and queue a fetch if not. + entitlementResource := grant.GetEntitlement().GetResource() + _, err := s.store.GetResource(ctx, &reader_v2.ResourcesReaderServiceGetResourceRequest{ + ResourceId: entitlementResource.GetId(), + }) + if err != nil { + if !errors.Is(err, sql.ErrNoRows) { + return err + } + + erId := entitlementResource.GetId() + prId := entitlementResource.GetParentResourceId() + resource, err := s.getResourceFromConnector(ctx, erId, prId) + if err != nil { + l.Error("error fetching entitlement resource", zap.Error(err)) + return err + } + if resource == nil { + continue + } + if err := s.store.PutResources(ctx, resource); err != nil { + return err + } + } + + principalResource := grant.GetPrincipal() + _, err = s.store.GetResource(ctx, &reader_v2.ResourcesReaderServiceGetResourceRequest{ + ResourceId: principalResource.GetId(), + }) + if err != nil { + if !errors.Is(err, sql.ErrNoRows) { + return err + } + + // Principal resource is not in the DB, so try to fetch it from the connector. + resource, err := s.getResourceFromConnector(ctx, principalResource.GetId(), principalResource.GetParentResourceId()) + if err != nil { + l.Error("error fetching principal resource", zap.Error(err)) + return err + } + if resource == nil { + continue + } + if err := s.store.PutResources(ctx, resource); err != nil { + return err + } + } } err = s.store.PutGrants(ctx, grants...) if err != nil { @@ -2623,6 +2719,18 @@ func WithTargetedSyncResourceIDs(resourceIDs []string) SyncOpt { } } +func WithOnlyExpandGrants() SyncOpt { + return func(s *syncer) { + s.onlyExpandGrants = true + } +} + +func WithSyncID(syncID string) SyncOpt { + return func(s *syncer) { + s.syncID = syncID + } +} + // NewSyncer returns a new syncer object. func NewSyncer(ctx context.Context, c types.ConnectorClient, opts ...SyncOpt) (Syncer, error) { s := &syncer{ diff --git a/vendor/github.com/conductorone/baton-sdk/pkg/synccompactor/compactor.go b/vendor/github.com/conductorone/baton-sdk/pkg/synccompactor/compactor.go index db442899..535c8ffe 100644 --- a/vendor/github.com/conductorone/baton-sdk/pkg/synccompactor/compactor.go +++ b/vendor/github.com/conductorone/baton-sdk/pkg/synccompactor/compactor.go @@ -7,12 +7,13 @@ import ( "io" "os" "path" - "runtime" - "syscall" + "path/filepath" reader_v2 "github.com/conductorone/baton-sdk/pb/c1/reader/v2" "github.com/conductorone/baton-sdk/pkg/dotc1z" c1zmanager "github.com/conductorone/baton-sdk/pkg/dotc1z/manager" + "github.com/conductorone/baton-sdk/pkg/sdk" + "github.com/conductorone/baton-sdk/pkg/sync" sync_compactor "github.com/conductorone/baton-sdk/pkg/synccompactor/naive" "github.com/grpc-ecosystem/go-grpc-middleware/logging/zap/ctxzap" "go.uber.org/zap" @@ -20,8 +21,9 @@ import ( type Compactor struct { entries []*CompactableSync - destDir string + tmpDir string + destDir string } type CompactableSync struct { @@ -33,39 +35,42 @@ var ErrNotEnoughFilesToCompact = errors.New("must provide two or more files to c type Option func(*Compactor) -// WithTmpDir sets the temporary directory for intermediate files during compaction. -func WithTmpDir(tmpDir string) Option { +// WithTmpDir sets the working directory where files will be created and edited during compaction. +// If not provided, the temporary directory will be used. +func WithTmpDir(tempDir string) Option { return func(c *Compactor) { - c.tmpDir = tmpDir + c.tmpDir = tempDir } } -func NewCompactor(ctx context.Context, destDir string, compactableSyncs []*CompactableSync, opts ...Option) (*Compactor, error) { +func NewCompactor(ctx context.Context, outputDir string, compactableSyncs []*CompactableSync, opts ...Option) (*Compactor, func() error, error) { if len(compactableSyncs) < 2 { - return nil, ErrNotEnoughFilesToCompact + return nil, nil, ErrNotEnoughFilesToCompact } - c := &Compactor{entries: compactableSyncs, destDir: destDir} + c := &Compactor{entries: compactableSyncs, destDir: outputDir} for _, opt := range opts { opt(c) } - return c, nil -} - -func removeIntermediateFiles(intermediates []string, preserveLast bool) error { - // The last one is our "base" so we don't want to remove that one - if preserveLast { - intermediates = intermediates[:len(intermediates)-1] + // If no tmpDir is provided, use the tmpDir + if c.tmpDir == "" { + c.tmpDir = os.TempDir() } - for _, intermediateFile := range intermediates { - err := os.Remove(intermediateFile) - // Weird case if the file doesn't exist but it's "fine". - if err != nil && !errors.Is(err, os.ErrNotExist) { + tmpDir, err := os.MkdirTemp(c.tmpDir, "baton-sync-compactor-") + if err != nil { + return nil, nil, err + } + c.tmpDir = tmpDir + + cleanup := func() error { + if err := os.RemoveAll(c.tmpDir); err != nil { return err } + return nil } - return nil + + return c, cleanup, nil } func (c *Compactor) Compact(ctx context.Context) (*CompactableSync, error) { @@ -73,63 +78,72 @@ func (c *Compactor) Compact(ctx context.Context) (*CompactableSync, error) { return nil, nil } - intermediates := make([]string, 0, len(c.entries)-1) - base := c.entries[0] for i := 1; i < len(c.entries); i++ { applied := c.entries[i] compactable, err := c.doOneCompaction(ctx, base, applied) if err != nil { - if err := removeIntermediateFiles(intermediates, false); err != nil { - return nil, err - } return nil, err } - // Collect all the intermediate files we create to remove at the end - intermediates = append(intermediates, compactable.FilePath) + base = compactable } - if len(intermediates) > 0 { - if err := removeIntermediateFiles(intermediates, true); err != nil { - return nil, err - } + l := ctxzap.Extract(ctx) + // Grant expansion doesn't use the connector interface at all, so giving syncer an empty connector is safe... for now. + // If that ever changes, we should implement a file connector that is a wrapper around the reader. + emptyConnector, err := sdk.NewEmptyConnector() + if err != nil { + l.Error("error creating empty connector", zap.Error(err)) + return nil, err + } + + // Use syncer to expand grants. + // TODO: Handle external resources. + syncer, err := sync.NewSyncer( + ctx, + emptyConnector, + sync.WithC1ZPath(base.FilePath), + sync.WithSyncID(base.SyncID), + sync.WithOnlyExpandGrants(), + ) + if err != nil { + l.Error("error creating syncer", zap.Error(err)) + return nil, err + } + + if err := syncer.Sync(ctx); err != nil { + l.Error("error syncing with grant expansion", zap.Error(err)) + return nil, err + } + if err := syncer.Close(ctx); err != nil { + l.Error("error closing syncer", zap.Error(err)) + return nil, err } // Move last compacted file to the destination dir finalPath := path.Join(c.destDir, fmt.Sprintf("compacted-%s.c1z", base.SyncID)) - // Attempt to move via rename - if err := os.Rename(base.FilePath, finalPath); err != nil { - var linkErr *os.LinkError - if errors.As(err, &linkErr) { - // Mac err table: https://developer.apple.com/library/archive/documentation/System/Conceptual/ManPages_iPhoneOS/man2/intro.2.html - // Win err table: https://learn.microsoft.com/en-us/openspecs/windows_protocols/ms-erref/18d8fbe8-a967-4f1c-ae50-99ca8e491d2d?redirectedfrom=MSDN - if errors.Is(linkErr.Err, syscall.Errno(0x12)) || (runtime.GOOS == "windows" && errors.Is(linkErr.Err, syscall.Errno(0x11))) { - // If rename doesn't work do a full create/copy - if err := mvFile(base.FilePath, finalPath); err != nil { - // Return if mv file failed - return nil, err - } - } else { - // Return if it's a different kind of link err - return nil, err - } - } else { - // Return if it's not a link err + if err := cpFile(base.FilePath, finalPath); err != nil { + return nil, err + } + + if !filepath.IsAbs(finalPath) { + abs, err := filepath.Abs(finalPath) + if err != nil { return nil, err } + finalPath = abs } - base.FilePath = finalPath - - return base, nil + return &CompactableSync{FilePath: finalPath, SyncID: base.SyncID}, nil } -func mvFile(sourcePath string, destPath string) error { +func cpFile(sourcePath string, destPath string) error { source, err := os.Open(sourcePath) if err != nil { return fmt.Errorf("failed to open source file: %w", err) } + defer source.Close() destination, err := os.Create(destPath) if err != nil { @@ -142,16 +156,6 @@ func mvFile(sourcePath string, destPath string) error { return fmt.Errorf("failed to copy file: %w", err) } - // Explicitly close the source file before removing it - if err := source.Close(); err != nil { - return err - } - - err = os.Remove(sourcePath) - if err != nil { - return fmt.Errorf("failed to remove original file: %w", err) - } - return nil } @@ -194,17 +198,18 @@ func (c *Compactor) doOneCompaction(ctx context.Context, base *CompactableSync, zap.String("base_sync", base.SyncID), zap.String("applied_file", applied.FilePath), zap.String("applied_sync", applied.SyncID), + zap.String("tmp_dir", c.tmpDir), ) - filePath := fmt.Sprintf("compacted-%s-%s.c1z", base.SyncID, applied.SyncID) - - opts := []dotc1z.C1ZOption{dotc1z.WithPragma("journal_mode", "WAL")} - if c.tmpDir != "" { - opts = append(opts, dotc1z.WithTmpDir(c.tmpDir)) + opts := []dotc1z.C1ZOption{ + dotc1z.WithPragma("journal_mode", "WAL"), + dotc1z.WithTmpDir(c.tmpDir), } - newFile, err := dotc1z.NewC1ZFile(ctx, filePath, opts...) + fileName := fmt.Sprintf("compacted-%s-%s.c1z", base.SyncID, applied.SyncID) + newFile, err := dotc1z.NewC1ZFile(ctx, path.Join(c.tmpDir, fileName), opts...) if err != nil { + l.Error("doOneCompaction failed: could not create c1z file", zap.Error(err)) return nil, err } defer func() { _ = newFile.Close() }() diff --git a/vendor/github.com/conductorone/baton-sdk/pkg/tasks/local/compactor.go b/vendor/github.com/conductorone/baton-sdk/pkg/tasks/local/compactor.go index 34daa7a6..1153134f 100644 --- a/vendor/github.com/conductorone/baton-sdk/pkg/tasks/local/compactor.go +++ b/vendor/github.com/conductorone/baton-sdk/pkg/tasks/local/compactor.go @@ -44,10 +44,14 @@ func (m *localCompactor) Process(ctx context.Context, task *v1.Task, cc types.Co defer span.End() log := ctxzap.Extract(ctx) - compactor, err := synccompactor.NewCompactor(ctx, m.outputPath, m.compactableSyncs) + compactor, cleanup, err := synccompactor.NewCompactor(ctx, m.outputPath, m.compactableSyncs) if err != nil { return err } + defer func() { + _ = cleanup() + }() + compacted, err := compactor.Compact(ctx) if err != nil { return err diff --git a/vendor/modules.txt b/vendor/modules.txt index 4b5b4b99..8daf8c5c 100644 --- a/vendor/modules.txt +++ b/vendor/modules.txt @@ -162,7 +162,7 @@ github.com/benbjohnson/clock # github.com/cenkalti/backoff/v4 v4.3.0 ## explicit; go 1.18 github.com/cenkalti/backoff/v4 -# github.com/conductorone/baton-sdk v0.3.7 +# github.com/conductorone/baton-sdk v0.3.8 ## explicit; go 1.23.4 github.com/conductorone/baton-sdk/internal/connector github.com/conductorone/baton-sdk/pb/c1/c1z/v1