diff --git a/.github/workflows/ci-core.yml b/.github/workflows/ci-core.yml index c72ff3b8de3..7efb2a8e3cf 100644 --- a/.github/workflows/ci-core.yml +++ b/.github/workflows/ci-core.yml @@ -391,7 +391,7 @@ jobs: id: wasm-key shell: bash run: | - echo "key=${{ runner.os }}-${{ runner.arch }}-wasmcache-${{ hashFiles('go.mod', 'go.sum', 'core/internal/testutils/wasmtest/**', 'core/capabilities/compute/**', 'core/services/workflows/**') }}" >> "$GITHUB_OUTPUT" + echo "key=${{ runner.os }}-${{ runner.arch }}-wasmcache-${{ hashFiles('go.mod', 'go.sum', 'core/internal/testutils/wasmtest/**', 'core/services/workflows/**') }}" >> "$GITHUB_OUTPUT" - name: Restore WASM test binaries cache if: ${{ matrix.type.should-run == 'true' }} diff --git a/.gitignore b/.gitignore index 2cf173ed28d..b8eeb296d02 100644 --- a/.gitignore +++ b/.gitignore @@ -52,13 +52,11 @@ cl_backup_*.tar.gz # Test artifacts core/cmd/TestClient_ImportExportP2PKeyBundle_test_key.json -race.* -*output.txt -/golangci-lint/ -.covdata -core/services/job/testdata/wasm/testmodule.wasm -core/services/job/testdata/wasm/testmodule.br -temp-repo + race.* + *output.txt + /golangci-lint/ + .covdata + temp-repo diagnose-*/ diagnose-attempted-fixes-*.jsonl diff --git a/core/services/feeds/service.go b/core/services/feeds/service.go index 59c5ba2514c..b2dc4c8dab3 100644 --- a/core/services/feeds/service.go +++ b/core/services/feeds/service.go @@ -623,18 +623,11 @@ func (s *service) DeleteJob(ctx context.Context, args *DeleteJobArgs) (int64, er return 0, errors.New("cannot delete a job proposal belonging to another feeds manager") } - // Try to delete as workflow job first (with auto-cancellation), fallback to simple deletion if not applicable - deleted, err := s.tryDeleteWithWorkflowCancellation(ctx, proposal, logger) - if err != nil { - return 0, err - } - - if !deleted { - // For non-workflow jobs: simple proposal deletion (no cancellation, just job_proposal delete) - if err = s.deleteSimpleJobProposal(ctx, proposal, logger); err != nil { - return 0, err - } + if err = s.orm.DeleteProposal(ctx, proposal.ID); err != nil { + logger.Errorw("Failed to delete the proposal", "err", err) + return 0, fmt.Errorf("DeleteProposal failed: %w", err) } + logger.Infow("Successfully deleted job proposal", "jobProposalID", proposal.ID) if err = s.observeJobProposalCounts(ctx); err != nil { logger.Errorw("Failed to push metrics for job proposal deletion", "err", err) @@ -848,7 +841,7 @@ func (s *service) ProposeJob(ctx context.Context, args *ProposeJobArgs) (int64, switch { case err != nil: logger.Errorw("Failed to validate spec while checking for workflow", "err", err) - case slices.Contains([]job.Type{job.Workflow, job.CRESettings}, jobType): + case jobType == job.CRESettings: promWorkflowRequests.Inc() promFeedsWorkflowRequests.Inc() err = s.ApproveSpec(ctx, specID, true) @@ -1035,15 +1028,6 @@ func (s *service) ApproveSpec(ctx context.Context, id int64, force bool) error { return errors.Wrap(txerr, "FindOCR2JobIDByAddress failed") } } - case job.Workflow: - existingJobID, txerr = tx.jobORM.FindJobIDByWorkflow(ctx, *j.WorkflowSpec) - if txerr != nil { - // Return an error if the repository errors. If there is a not found - // error we want to continue with approving the job. - if !errors.Is(txerr, sql.ErrNoRows) { - return fmt.Errorf("failed while checking for existing workflow job: %w", txerr) - } - } case job.CCIP: existingJobID, txerr = tx.jobORM.FindJobIDByCapabilityNameAndVersion(ctx, *j.CCIPSpec) // Return an error if the repository errors. If there is a not found @@ -1900,101 +1884,3 @@ func (ns NullService) UpdateSpecDefinition(ctx context.Context, id int64, spec s func (ns NullService) Unsafe_SetConnectionsManager(_ ConnectionsManager) {} //revive:enable - -// deleteSimpleJobProposal deletes a simple (non-workflow) job proposal. -// This only removes the proposal without any cancellation, unlike workflow jobs -func (s *service) deleteSimpleJobProposal(ctx context.Context, proposal *JobProposal, logger logger.Logger) error { - if err := s.orm.DeleteProposal(ctx, proposal.ID); err != nil { - logger.Errorw("Failed to delete the proposal", "err", err) - return fmt.Errorf("DeleteProposal failed: %w", err) - } - - logger.Infow("Successfully deleted simple job proposal", "jobProposalID", proposal.ID) - return nil -} - -// tryDeleteWithWorkflowCancellation attempts to delete a job as a workflow job. -// Returns true if the job was successfully deleted as a workflow, false if it's not a workflow job. -// Returns an error if deletion failed. -func (s *service) tryDeleteWithWorkflowCancellation(ctx context.Context, proposal *JobProposal, logger logger.Logger) (bool, error) { - // Early return if no external job ID (we won't find a job to delete without it) - if !proposal.ExternalJobID.Valid { - logger.Debugw("Proposal has no ExternalJobID, skipping workflow job deletion", "proposalID", proposal.ID) - return false, nil - } - - // Try to find the job by external job ID - jobFound, err := s.jobORM.FindJobByExternalJobID(ctx, proposal.ExternalJobID.UUID) - if err != nil { - logger.Warnw("Failed to find job by external job ID, skipping workflow job deletion", - "externalJobID", proposal.ExternalJobID.UUID, "err", err) - return false, nil - } - - // Check if this is actually a workflow job - if jobFound.WorkflowSpecID == nil { - logger.Debugw("Job is not a workflow job, skipping workflow job deletion", - "jobID", jobFound.ID, "jobType", jobFound.Type) - return false, nil - } - - // Get the approved spec for workflow cancellation - jpSpec, err := s.orm.GetApprovedSpec(ctx, proposal.ID) - if err != nil { - logger.Errorw("GetApprovedSpec failed - cannot proceed with workflow job deletion", - "proposalID", proposal.ID, "err", err, "jobName", jobFound.Name) - return false, nil - } - - // All validations passed - proceed with workflow job deletion - logger.Debugw("Proceeding with workflow job deletion", - "proposalID", proposal.ID, "jobID", jobFound.ID, "specID", jpSpec.ID) - - return true, s.deleteWorkflowJobWithTransaction(ctx, *proposal, jobFound, *jpSpec, logger) -} - -// deleteWorkflowJobWithTransaction performs workflow job deletion with auto-cancellation within a transaction. -func (s *service) deleteWorkflowJobWithTransaction(ctx context.Context, proposal JobProposal, job job.Job, jpSpec JobProposalSpec, logger logger.Logger) error { - if job.WorkflowSpecID == nil { - return errors.New("job WorkflowSpecID is nil, cannot delete workflow job") - } - jobSpecID := int64(*job.WorkflowSpecID) - - fmsClient, err := s.connMgr.GetClient(proposal.FeedsManagerID) - if err != nil { - logger.Errorw("Failed to get FMS client", "jobProposalID", proposal.ID, "jobProposalSpecID", jpSpec.ID, "err", err, "name", job.Name) - return fmt.Errorf("failed to get FMS client for workflow spec cancellation: %w", err) - } - - cancelLogger := logger.With("job_proposal_spec_id", jpSpec.ID, "jobSpecID", jobSpecID) - - err = s.transact(ctx, func(tx datasources) error { - if txerr := tx.orm.DeleteProposal(ctx, proposal.ID); txerr != nil { - return fmt.Errorf("DeleteProposal failed: %w", txerr) - } - - if txerr := tx.orm.CancelSpec(ctx, jpSpec.ID); txerr != nil { - return txerr - } - - if serr := s.jobSpawner.DeleteJob(ctx, tx.ds, job.ID); serr != nil { - return fmt.Errorf("DeleteJob failed: %w", serr) - } - - if _, err = fmsClient.CancelledJob(ctx, &pb.CancelledJobRequest{ - Uuid: proposal.RemoteUUID.String(), - Version: int64(jpSpec.Version), - }); err != nil { - return err - } - - return nil - }) - if err != nil { - cancelLogger.Errorw("Failed to auto-cancel workflow spec", "err", err, "name", job.Name) - return fmt.Errorf("failed to auto-cancel workflow spec (job proposal spec ID: %d): %w", jpSpec.ID, err) - } - - logger.Infow("Successfully auto-cancelled a workflow spec", "jobProposalID", proposal.ID, "jobProposalSpecID", jpSpec.ID, "jobSpecID", jobSpecID, "name", job.Name) - return nil -} diff --git a/core/services/feeds/service_test.go b/core/services/feeds/service_test.go index 9ab9a580497..1df02b60f2c 100644 --- a/core/services/feeds/service_test.go +++ b/core/services/feeds/service_test.go @@ -1133,18 +1133,6 @@ func Test_Service_DeleteJob(t *testing.T) { Status: feeds.JobProposalStatusApproved, } - wfSpecID = int32(4321) - workflowJob = job.Job{ - ID: 1, - WorkflowSpecID: &wfSpecID, - } - jobProposalSpec = &feeds.JobProposalSpec{ - ID: 20, - Status: feeds.SpecStatusApproved, - JobProposalID: approved.ID, - Version: 1, - } - httpTimeout = *commonconfig.MustNewDuration(1 * time.Second) ) @@ -1161,7 +1149,6 @@ func Test_Service_DeleteJob(t *testing.T) { svc.orm.On("GetJobProposalByRemoteUUID", mock.Anything, approved.RemoteUUID).Return(&approved, nil) svc.orm.On("DeleteProposal", mock.Anything, approved.ID).Return(nil) svc.orm.On("CountJobProposalsByStatus", mock.Anything).Return(&feeds.JobProposalCounts{}, nil) - svc.jobORM.On("FindJobByExternalJobID", mock.Anything, approved.ExternalJobID.UUID).Return(job.Job{}, sql.ErrNoRows) }, args: args, wantID: approved.ID, @@ -1200,187 +1187,11 @@ func Test_Service_DeleteJob(t *testing.T) { name: "Delete proposal error", before: func(svc *TestService) { svc.orm.On("GetJobProposalByRemoteUUID", mock.Anything, approved.RemoteUUID).Return(&approved, nil) - svc.jobORM.On("FindJobByExternalJobID", mock.Anything, approved.ExternalJobID.UUID).Return(job.Job{}, sql.ErrNoRows) svc.orm.On("DeleteProposal", mock.Anything, approved.ID).Return(errors.New("orm error")) }, args: args, wantErr: "DeleteProposal failed", }, - { - name: "Delete workflow-spec with auto-cancellation", - before: func(svc *TestService) { - svc.orm.On("GetJobProposalByRemoteUUID", mock.Anything, approved.RemoteUUID).Return(&approved, nil) - svc.orm.On("DeleteProposal", mock.Anything, approved.ID).Return(nil) - svc.orm.On("CountJobProposalsByStatus", mock.Anything).Return(&feeds.JobProposalCounts{}, nil) - svc.jobORM.On("FindJobByExternalJobID", mock.Anything, approved.ExternalJobID.UUID).Return(workflowJob, nil) - svc.orm.On("GetApprovedSpec", mock.Anything, approved.ID).Return(jobProposalSpec, nil) - - svc.connMgr.On("GetClient", mock.Anything).Return(svc.fmsClient, nil) - - svc.orm.On("CancelSpec", mock.Anything, jobProposalSpec.ID).Return(nil) - svc.spawner.On("DeleteJob", mock.Anything, mock.Anything, workflowJob.ID).Return(nil) - - svc.fmsClient.On("CancelledJob", - mock.MatchedBy(func(ctx context.Context) bool { return true }), - &proto.CancelledJobRequest{ - Uuid: approved.RemoteUUID.String(), - Version: int64(jobProposalSpec.Version), - }, - ).Return(&proto.CancelledJobResponse{}, nil) - svc.orm.On("CountJobProposalsByStatus", mock.Anything).Return(&feeds.JobProposalCounts{}, nil) - svc.orm.On("WithDataSource", mock.Anything).Return(feeds.ORM(svc.orm)) - svc.jobORM.On("WithDataSource", mock.Anything).Return(job.ORM(svc.jobORM)) - }, - args: args, - wantID: approved.ID, - }, - { - name: "Delete workflow-spec transaction rollback on FMS client error", - before: func(svc *TestService) { - svc.orm.On("GetJobProposalByRemoteUUID", mock.Anything, approved.RemoteUUID).Return(&approved, nil) - svc.jobORM.On("FindJobByExternalJobID", mock.Anything, approved.ExternalJobID.UUID).Return(workflowJob, nil) - svc.orm.On("GetApprovedSpec", mock.Anything, approved.ID).Return(jobProposalSpec, nil) - - svc.connMgr.On("GetClient", mock.Anything).Return(svc.fmsClient, nil) - - // These should be called but then rolled back due to FMS error - svc.orm.On("DeleteProposal", mock.Anything, approved.ID).Return(nil) - svc.orm.On("CancelSpec", mock.Anything, jobProposalSpec.ID).Return(nil) - svc.spawner.On("DeleteJob", mock.Anything, mock.Anything, workflowJob.ID).Return(nil) - - // FMS client call fails - this should cause transaction rollback - svc.fmsClient.On("CancelledJob", - mock.MatchedBy(func(ctx context.Context) bool { return true }), - &proto.CancelledJobRequest{ - Uuid: approved.RemoteUUID.String(), - Version: int64(jobProposalSpec.Version), - }, - ).Return(nil, errors.New("FMS client timeout")) - - svc.orm.On("WithDataSource", mock.Anything).Return(feeds.ORM(svc.orm)) - svc.jobORM.On("WithDataSource", mock.Anything).Return(job.ORM(svc.jobORM)) - }, - args: args, - wantErr: "failed to auto-cancel workflow spec", - }, - { - name: "Delete workflow-spec transaction rollback on job deletion error", - before: func(svc *TestService) { - svc.orm.On("GetJobProposalByRemoteUUID", mock.Anything, approved.RemoteUUID).Return(&approved, nil) - svc.jobORM.On("FindJobByExternalJobID", mock.Anything, approved.ExternalJobID.UUID).Return(workflowJob, nil) - svc.orm.On("GetApprovedSpec", mock.Anything, approved.ID).Return(jobProposalSpec, nil) - - svc.connMgr.On("GetClient", mock.Anything).Return(svc.fmsClient, nil) - - // These should be called but then rolled back due to job deletion error - svc.orm.On("DeleteProposal", mock.Anything, approved.ID).Return(nil) - svc.orm.On("CancelSpec", mock.Anything, jobProposalSpec.ID).Return(nil) - - // Job deletion fails - this should cause transaction rollback - svc.spawner.On("DeleteJob", mock.Anything, mock.Anything, workflowJob.ID).Return(errors.New("job deletion failed")) - - svc.orm.On("WithDataSource", mock.Anything).Return(feeds.ORM(svc.orm)) - svc.jobORM.On("WithDataSource", mock.Anything).Return(job.ORM(svc.jobORM)) - }, - args: args, - wantErr: "failed to auto-cancel workflow spec", - }, - { - name: "GetClient error for workflow cancellation", - before: func(svc *TestService) { - svc.orm.On("GetJobProposalByRemoteUUID", mock.Anything, approved.RemoteUUID).Return(&approved, nil) - svc.jobORM.On("FindJobByExternalJobID", mock.Anything, approved.ExternalJobID.UUID).Return(workflowJob, nil) - svc.orm.On("GetApprovedSpec", mock.Anything, approved.ID).Return(jobProposalSpec, nil) - - svc.connMgr.On("GetClient", mock.Anything).Return(nil, errors.New("connection manager error")) - }, - args: args, - wantErr: "failed to get FMS client for workflow spec cancellation", - }, - { - name: "GetApprovedSpec error for workflow job - fallback to simple deletion", - before: func(svc *TestService) { - svc.orm.On("GetJobProposalByRemoteUUID", mock.Anything, approved.RemoteUUID).Return(&approved, nil) - svc.jobORM.On("FindJobByExternalJobID", mock.Anything, approved.ExternalJobID.UUID).Return(workflowJob, nil) - svc.orm.On("GetApprovedSpec", mock.Anything, approved.ID).Return(nil, errors.New("no approved spec")) - - // Should fallback to simple proposal deletion - svc.orm.On("DeleteProposal", mock.Anything, approved.ID).Return(nil) - svc.orm.On("CountJobProposalsByStatus", mock.Anything).Return(&feeds.JobProposalCounts{}, nil) - }, - args: args, - wantID: approved.ID, - }, - { - name: "Proposal with ExternalJobID but job not found - simple deletion", - before: func(svc *TestService) { - proposalWithJobID := approved - proposalWithJobID.ExternalJobID = uuid.NullUUID{UUID: uuid.New(), Valid: true} - - svc.orm.On("GetJobProposalByRemoteUUID", mock.Anything, approved.RemoteUUID).Return(&proposalWithJobID, nil) - svc.jobORM.On("FindJobByExternalJobID", mock.Anything, proposalWithJobID.ExternalJobID.UUID).Return(job.Job{}, sql.ErrNoRows) - svc.orm.On("DeleteProposal", mock.Anything, approved.ID).Return(nil) - svc.orm.On("CountJobProposalsByStatus", mock.Anything).Return(&feeds.JobProposalCounts{}, nil) - }, - args: args, - wantID: approved.ID, - }, - { - name: "Proposal with ExternalJobID but job is not workflow type - simple deletion", - before: func(svc *TestService) { - proposalWithJobID := approved - proposalWithJobID.ExternalJobID = uuid.NullUUID{UUID: uuid.New(), Valid: true} - - nonWorkflowJob := job.Job{ - ID: 2, - WorkflowSpecID: nil, // Not a workflow job - } - - svc.orm.On("GetJobProposalByRemoteUUID", mock.Anything, approved.RemoteUUID).Return(&proposalWithJobID, nil) - svc.jobORM.On("FindJobByExternalJobID", mock.Anything, proposalWithJobID.ExternalJobID.UUID).Return(nonWorkflowJob, nil) - svc.orm.On("DeleteProposal", mock.Anything, approved.ID).Return(nil) - svc.orm.On("CountJobProposalsByStatus", mock.Anything).Return(&feeds.JobProposalCounts{}, nil) - }, - args: args, - wantID: approved.ID, - }, - { - name: "DeleteProposal error in workflow cancellation path", - before: func(svc *TestService) { - svc.orm.On("GetJobProposalByRemoteUUID", mock.Anything, approved.RemoteUUID).Return(&approved, nil) - svc.jobORM.On("FindJobByExternalJobID", mock.Anything, approved.ExternalJobID.UUID).Return(workflowJob, nil) - svc.orm.On("GetApprovedSpec", mock.Anything, approved.ID).Return(jobProposalSpec, nil) - - svc.connMgr.On("GetClient", mock.Anything).Return(svc.fmsClient, nil) - - // DeleteProposal fails in workflow cancellation path - svc.orm.On("DeleteProposal", mock.Anything, approved.ID).Return(errors.New("delete proposal failed")) - - svc.orm.On("WithDataSource", mock.Anything).Return(feeds.ORM(svc.orm)) - svc.jobORM.On("WithDataSource", mock.Anything).Return(job.ORM(svc.jobORM)) - }, - args: args, - wantErr: "failed to auto-cancel workflow spec", - }, - { - name: "CancelSpec error in workflow cancellation path", - before: func(svc *TestService) { - svc.orm.On("GetJobProposalByRemoteUUID", mock.Anything, approved.RemoteUUID).Return(&approved, nil) - svc.jobORM.On("FindJobByExternalJobID", mock.Anything, approved.ExternalJobID.UUID).Return(workflowJob, nil) - svc.orm.On("GetApprovedSpec", mock.Anything, approved.ID).Return(jobProposalSpec, nil) - - svc.connMgr.On("GetClient", mock.Anything).Return(svc.fmsClient, nil) - - svc.orm.On("DeleteProposal", mock.Anything, approved.ID).Return(nil) - // CancelSpec fails - svc.orm.On("CancelSpec", mock.Anything, jobProposalSpec.ID).Return(errors.New("cancel spec failed")) - - svc.orm.On("WithDataSource", mock.Anything).Return(feeds.ORM(svc.orm)) - svc.jobORM.On("WithDataSource", mock.Anything).Return(job.ORM(svc.jobORM)) - }, - args: args, - wantErr: "failed to auto-cancel workflow spec", - }, { name: "observeJobProposalCounts error - success with warning log", before: func(svc *TestService) { @@ -1388,7 +1199,6 @@ func Test_Service_DeleteJob(t *testing.T) { svc.orm.On("DeleteProposal", mock.Anything, approved.ID).Return(nil) // observeJobProposalCounts fails but shouldn't cause DeleteJob to fail svc.orm.On("CountJobProposalsByStatus", mock.Anything).Return(nil, errors.New("metrics error")) - svc.jobORM.On("FindJobByExternalJobID", mock.Anything, approved.ExternalJobID.UUID).Return(job.Job{}, sql.ErrNoRows) }, args: args, wantID: approved.ID, diff --git a/core/services/job/job_orm_test.go b/core/services/job/job_orm_test.go index 44816061fc4..63bdf9d1445 100644 --- a/core/services/job/job_orm_test.go +++ b/core/services/job/job_orm_test.go @@ -10,7 +10,6 @@ import ( "github.com/ethereum/go-ethereum/common" "github.com/google/uuid" "github.com/lib/pq" - "github.com/pelletier/go-toml/v2" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" "gopkg.in/guregu/null.v4" @@ -19,7 +18,6 @@ import ( "github.com/smartcontractkit/chainlink-common/pkg/sqlutil" "github.com/smartcontractkit/chainlink-common/pkg/types" "github.com/smartcontractkit/chainlink-common/pkg/utils/jsonserializable" - pkgworkflows "github.com/smartcontractkit/chainlink-common/pkg/workflows" "github.com/smartcontractkit/chainlink-evm/pkg/assets" configtoml "github.com/smartcontractkit/chainlink-evm/pkg/config/toml" "github.com/smartcontractkit/chainlink-evm/pkg/keys" @@ -48,7 +46,6 @@ import ( "github.com/smartcontractkit/chainlink/v2/core/services/standardcapabilities" "github.com/smartcontractkit/chainlink/v2/core/services/streams" "github.com/smartcontractkit/chainlink/v2/core/services/vrf/vrfcommon" - "github.com/smartcontractkit/chainlink/v2/core/services/workflows/artifacts" "github.com/smartcontractkit/chainlink/v2/core/testdata/testspecs" "github.com/smartcontractkit/chainlink/v2/core/utils/testutils/heavyweight" ) @@ -710,6 +707,18 @@ func TestORM_CreateJob_EVMChainID_Validation(t *testing.T) { assert.ErrorIs(t, err, job.ErrJobTypeRemoved) }) + t.Run("workflow job creation is rejected", func(t *testing.T) { + t.Parallel() + jb := job.Job{ + Type: job.Workflow, + WorkflowSpec: &job.WorkflowSpec{}, + SchemaVersion: 1, + } + err := jobORM.CreateJob(t.Context(), &jb) + require.ErrorIs(t, err, job.ErrJobTypeRemoved) + require.ErrorContains(t, err, `cannot create job of type "workflow"`) + }) + t.Run("evm chain id validation for vrf works", func(t *testing.T) { jb := job.Job{ Type: job.VRF, @@ -1882,215 +1891,6 @@ func Test_CountPipelineRunsByJobID(t *testing.T) { }) } -func Test_ORM_FindJobByWorkflow(t *testing.T) { - addr1 := "0x0123456789012345678901234567890123456789" - addr2 := "0xabcdefabcdefabcdefabcdefabcdefabcdefabcd" - t.Parallel() - type fields struct { - ds sqlutil.DataSource - } - type args struct { - spec *job.WorkflowSpec - before func(t *testing.T, o job.ORM, s *job.WorkflowSpec) int32 - } - tests := []struct { - name string - fields fields - args args - wantErr bool - }{ - { - name: "wf not job found", - fields: fields{ - ds: pgtest.NewSqlxDB(t), - }, - args: args{ - // before is nil, so no job is inserted - spec: &job.WorkflowSpec{ - ID: 1, - Workflow: pkgworkflows.WFYamlSpec(t, "workflow00", addr1), - }, - }, - wantErr: true, - }, - - { - name: "wf job found", - fields: fields{ - ds: pgtest.NewSqlxDB(t), - }, - args: args{ - spec: &job.WorkflowSpec{ - ID: 1, - Workflow: pkgworkflows.WFYamlSpec(t, "workflow01", addr1), - SpecType: job.YamlSpec, - }, - before: mustInsertWFJob, - }, - wantErr: false, - }, - - { - name: "wf wrong name", - fields: fields{ - ds: pgtest.NewSqlxDB(t), - }, - args: args{ - spec: &job.WorkflowSpec{ - ID: 1, - Workflow: pkgworkflows.WFYamlSpec(t, "workflow02", addr1), - }, - before: func(t *testing.T, o job.ORM, s *job.WorkflowSpec) int32 { - var c job.WorkflowSpec - c.ID = s.ID - c.Workflow = pkgworkflows.WFYamlSpec(t, "workflow99", addr1) // insert with mismatched name - c.SpecType = job.YamlSpec - c.SecretsID = s.SecretsID - return mustInsertWFJob(t, o, &c) - }, - }, - wantErr: true, - }, - { - name: "wf wrong owner", - fields: fields{ - ds: pgtest.NewSqlxDB(t), - }, - args: args{ - spec: &job.WorkflowSpec{ - ID: 1, - Workflow: pkgworkflows.WFYamlSpec(t, "workflow03", addr1), - }, - before: func(t *testing.T, o job.ORM, s *job.WorkflowSpec) int32 { - var c job.WorkflowSpec - c.ID = s.ID - c.Workflow = pkgworkflows.WFYamlSpec(t, "workflow03", addr2) // insert with mismatched owner - c.SecretsID = s.SecretsID - return mustInsertWFJob(t, o, &c) - }, - }, - wantErr: true, - }, - } - - for i, tt := range tests { - t.Run(tt.name, func(t *testing.T) { - ctx := t.Context() - ks := cltest.NewKeyStore(t, tt.fields.ds) - - secretsORM := artifacts.NewWorkflowRegistryDS(tt.fields.ds, logger.TestLogger(t)) - - sid, err := secretsORM.Create(ctx, "some-url.com", fmt.Sprintf("some-hash-%d", i), "some-contentz") - require.NoError(t, err) - tt.args.spec.SecretsID = sql.NullInt64{Int64: sid, Valid: true} - - pipelineORM := pipeline.NewORM(tt.fields.ds, logger.TestLogger(t), configtest.NewTestGeneralConfig(t).JobPipeline().MaxSuccessfulRuns()) - bridgesORM := bridges.NewORM(tt.fields.ds) - o := NewTestORM(t, tt.fields.ds, pipelineORM, bridgesORM, ks) - - var wantJobID int32 - if tt.args.before != nil { - wantJobID = tt.args.before(t, o, tt.args.spec) - } - - gotJ, err := o.FindJobIDByWorkflow(ctx, *tt.args.spec) - if (err != nil) != tt.wantErr { - t.Errorf("orm.FindJobByWorkflow() error = %v, wantErr %v", err, tt.wantErr) - return - } - - if err == nil { - assert.Equal(t, wantJobID, gotJ, "mismatch job id") - } - }) - } -} - -func Test_ORM_FindJobByWorkflow_Multiple(t *testing.T) { - addr1 := "0x012345678901234567890123456789012345ffff" - addr2 := "0xabcdefabcdefabcdefabcdefabcdefabcdef0000" - t.Parallel() - t.Run("multiple jobs", func(t *testing.T) { - db := pgtest.NewSqlxDB(t) - o := NewTestORM(t, - db, - pipeline.NewORM(db, - logger.TestLogger(t), - configtest.NewTestGeneralConfig(t).JobPipeline().MaxSuccessfulRuns()), - bridges.NewORM(db), - cltest.NewKeyStore(t, db)) - ctx := t.Context() - secretsORM := artifacts.NewWorkflowRegistryDS(db, logger.TestLogger(t)) - - sids := make([]int64, 0, 3) - for i := range 3 { - sid, err := secretsORM.Create(ctx, "some-url.com", fmt.Sprintf("some-hash-%d", i), "some-contentz") - require.NoError(t, err) - sids = append(sids, sid) - } - - wfYaml1 := pkgworkflows.WFYamlSpec(t, "workflow00", addr1) - s1 := job.WorkflowSpec{ - Workflow: wfYaml1, - SpecType: job.YamlSpec, - SecretsID: sql.NullInt64{Int64: sids[0], Valid: true}, - } - wantJobID1 := mustInsertWFJob(t, o, &s1) - - wfYaml2 := pkgworkflows.WFYamlSpec(t, "workflow01", addr1) - s2 := job.WorkflowSpec{ - Workflow: wfYaml2, - SpecType: job.YamlSpec, - SecretsID: sql.NullInt64{Int64: sids[1], Valid: true}, - } - wantJobID2 := mustInsertWFJob(t, o, &s2) - - wfYaml3 := pkgworkflows.WFYamlSpec(t, "workflow00", addr2) - s3 := job.WorkflowSpec{ - Workflow: wfYaml3, - SpecType: job.YamlSpec, - SecretsID: sql.NullInt64{Int64: sids[2], Valid: true}, - } - wantJobID3 := mustInsertWFJob(t, o, &s3) - - expectedIDs := []int32{wantJobID1, wantJobID2, wantJobID3} - for i, s := range []job.WorkflowSpec{s1, s2, s3} { - gotJ, err := o.FindJobIDByWorkflow(ctx, s) - require.NoError(t, err) - assert.Equal(t, expectedIDs[i], gotJ, "mismatch job id case %d, spec %v", i, s) - j, err := o.FindJob(ctx, expectedIDs[i]) - require.NoError(t, err) - assert.NotNil(t, j) - t.Logf("found job %v", j) - assert.Equal(t, j.WorkflowSpec.Workflow, s.Workflow) - assert.Equal(t, j.WorkflowSpec.WorkflowID, s.WorkflowID) - assert.Equal(t, j.WorkflowSpec.WorkflowOwner, s.WorkflowOwner) - assert.Equal(t, j.WorkflowSpec.WorkflowName, s.WorkflowName) - assert.Equal(t, job.YamlSpec, j.WorkflowSpec.SpecType) - } - }) -} - -func mustInsertWFJob(t *testing.T, orm job.ORM, s *job.WorkflowSpec) int32 { - t.Helper() - err := s.Validate(t.Context()) - require.NoError(t, err, "failed to validate spec %v", s) - ctx := t.Context() - _, err = toml.Marshal(s) - require.NoError(t, err, "failed to TOML marshal workflow spec %v", s) - j := job.Job{ - Type: job.Workflow, - WorkflowSpec: s, - ExternalJobID: uuid.New(), - Name: null.StringFrom(s.WorkflowOwner + "_" + s.WorkflowName), - SchemaVersion: 1, - } - - err = orm.CreateJob(ctx, &j) - require.NoError(t, err, "failed to insert job with wf spec %+v %s", s, err) - return j.ID -} - func mustInsertPipelineRun(t *testing.T, orm pipeline.ORM, j job.Job) pipeline.Run { t.Helper() ctx := t.Context() diff --git a/core/services/job/mocks/orm.go b/core/services/job/mocks/orm.go index 7c9b126da18..3feae06145d 100644 --- a/core/services/job/mocks/orm.go +++ b/core/services/job/mocks/orm.go @@ -657,63 +657,6 @@ func (_c *ORM_FindJobIDByCapabilityNameAndVersion_Call) RunAndReturn(run func(co return _c } -// FindJobIDByWorkflow provides a mock function with given fields: ctx, spec -func (_m *ORM) FindJobIDByWorkflow(ctx context.Context, spec job.WorkflowSpec) (int32, error) { - ret := _m.Called(ctx, spec) - - if len(ret) == 0 { - panic("no return value specified for FindJobIDByWorkflow") - } - - var r0 int32 - var r1 error - if rf, ok := ret.Get(0).(func(context.Context, job.WorkflowSpec) (int32, error)); ok { - return rf(ctx, spec) - } - if rf, ok := ret.Get(0).(func(context.Context, job.WorkflowSpec) int32); ok { - r0 = rf(ctx, spec) - } else { - r0 = ret.Get(0).(int32) - } - - if rf, ok := ret.Get(1).(func(context.Context, job.WorkflowSpec) error); ok { - r1 = rf(ctx, spec) - } else { - r1 = ret.Error(1) - } - - return r0, r1 -} - -// ORM_FindJobIDByWorkflow_Call is a *mock.Call that shadows Run/Return methods with type explicit version for method 'FindJobIDByWorkflow' -type ORM_FindJobIDByWorkflow_Call struct { - *mock.Call -} - -// FindJobIDByWorkflow is a helper method to define mock.On call -// - ctx context.Context -// - spec job.WorkflowSpec -func (_e *ORM_Expecter) FindJobIDByWorkflow(ctx interface{}, spec interface{}) *ORM_FindJobIDByWorkflow_Call { - return &ORM_FindJobIDByWorkflow_Call{Call: _e.mock.On("FindJobIDByWorkflow", ctx, spec)} -} - -func (_c *ORM_FindJobIDByWorkflow_Call) Run(run func(ctx context.Context, spec job.WorkflowSpec)) *ORM_FindJobIDByWorkflow_Call { - _c.Call.Run(func(args mock.Arguments) { - run(args[0].(context.Context), args[1].(job.WorkflowSpec)) - }) - return _c -} - -func (_c *ORM_FindJobIDByWorkflow_Call) Return(_a0 int32, _a1 error) *ORM_FindJobIDByWorkflow_Call { - _c.Call.Return(_a0, _a1) - return _c -} - -func (_c *ORM_FindJobIDByWorkflow_Call) RunAndReturn(run func(context.Context, job.WorkflowSpec) (int32, error)) *ORM_FindJobIDByWorkflow_Call { - _c.Call.Return(run) - return _c -} - // FindJobIDsWithBridge provides a mock function with given fields: ctx, name func (_m *ORM) FindJobIDsWithBridge(ctx context.Context, name string) ([]int32, error) { ret := _m.Called(ctx, name) diff --git a/core/services/job/models.go b/core/services/job/models.go index afaf1fc430f..7bf49ea202f 100644 --- a/core/services/job/models.go +++ b/core/services/job/models.go @@ -1,7 +1,6 @@ package job import ( - "context" "database/sql" "database/sql/driver" "encoding/json" @@ -22,7 +21,6 @@ import ( "github.com/smartcontractkit/chainlink-common/pkg/sqlutil" "github.com/smartcontractkit/chainlink-common/pkg/types" clnull "github.com/smartcontractkit/chainlink-common/pkg/utils/null" - "github.com/smartcontractkit/chainlink-common/pkg/workflows/sdk" "github.com/smartcontractkit/chainlink-evm/pkg/assets" "github.com/smartcontractkit/chainlink-evm/pkg/config/toml" evmtypes "github.com/smartcontractkit/chainlink-evm/pkg/types" @@ -810,9 +808,7 @@ type LiquidityBalancerSpec struct { type WorkflowSpecType string const ( - YamlSpec WorkflowSpecType = "yaml" - WASMFile WorkflowSpecType = "wasm_file" - DefaultSpecType = "" + WASMFile WorkflowSpecType = "wasm_file" ) type WorkflowSpecStatus string @@ -848,99 +844,6 @@ type WorkflowSpec struct { // StorageBytes is the workflow + config size in bytes. Set at registration // and not cleared by pausing the workflow StorageBytes int64 `toml:"-" db:"storage_bytes"` - - sdkWorkflow *sdk.WorkflowSpec - rawSpec []byte - config []byte -} - -var ( - ErrInvalidWorkflowID = errors.New("invalid workflow id") - ErrInvalidWorkflowYAMLSpec = errors.New("invalid workflow yaml spec") -) - -const ( - workflowIDLen = 64 // sha256 hash -) - -// Validate checks the workflow spec for correctness -func (w *WorkflowSpec) Validate(ctx context.Context) error { - s, err := w.SDKSpec(ctx) - if err != nil { - return err - } - - // For yaml-based workflow specs, use the owner & name fields defined there. - // For wasm workflows, use the `workflow_name` & `workflow_owner` fields directly from the job spec. - if s.Owner+s.Name != "" { - w.WorkflowOwner = strings.TrimPrefix(s.Owner, "0x") // the json schema validation ensures it is a hex string with 0x prefix, but the database does not store the prefix - w.WorkflowName = s.Name - } else { - w.WorkflowOwner = strings.TrimPrefix(w.WorkflowOwner, "0x") - } - - if len(w.WorkflowID) != workflowIDLen { - return fmt.Errorf("%w: incorrect length for id %s: expected %d, got %d", ErrInvalidWorkflowID, w.WorkflowID, workflowIDLen, len(w.WorkflowID)) - } - - return nil -} - -func (w *WorkflowSpec) SDKSpec(ctx context.Context) (sdk.WorkflowSpec, error) { - if w.sdkWorkflow != nil { - return *w.sdkWorkflow, nil - } - - workflowSpecFactory, ok := workflowSpecFactories[w.SpecType] - if !ok { - return sdk.WorkflowSpec{}, fmt.Errorf("unknown spec type %s", w.SpecType) - } - spec, rawSpec, cid, err := workflowSpecFactory.Spec(ctx, w.Workflow, w.Config) - if err != nil { - return sdk.WorkflowSpec{}, fmt.Errorf("spec factory failed: %w", err) - } - w.sdkWorkflow = &spec - w.rawSpec = rawSpec - w.WorkflowID = cid - return spec, nil -} - -func (w *WorkflowSpec) RawSpec(ctx context.Context) ([]byte, error) { - if w.rawSpec != nil { - return w.rawSpec, nil - } - - workflowSpecFactory, ok := workflowSpecFactories[w.SpecType] - if !ok { - return nil, fmt.Errorf("unknown spec type %s", w.SpecType) - } - - rs, err := workflowSpecFactory.RawSpec(ctx, w.Workflow, w.Config) - if err != nil { - return nil, err - } - - w.rawSpec = rs - return rs, nil -} - -func (w *WorkflowSpec) GetConfig(ctx context.Context) ([]byte, error) { - if w.config != nil { - return w.config, nil - } - - workflowSpecFactory, ok := workflowSpecFactories[w.SpecType] - if !ok { - return nil, fmt.Errorf("unknown spec type %s", w.SpecType) - } - - rs, err := workflowSpecFactory.Config(ctx, w.Config) - if err != nil { - return nil, err - } - - w.config = rs - return rs, nil } type StandardCapabilitiesConfig struct { diff --git a/core/services/job/models_test.go b/core/services/job/models_test.go index f721a1bf568..7f0ca25db7a 100644 --- a/core/services/job/models_test.go +++ b/core/services/job/models_test.go @@ -2,7 +2,6 @@ package job_test import ( _ "embed" - "fmt" "reflect" "testing" "time" @@ -15,7 +14,6 @@ import ( "github.com/smartcontractkit/chainlink-common/pkg/codec" "github.com/smartcontractkit/chainlink-common/pkg/sqlutil" "github.com/smartcontractkit/chainlink-common/pkg/types" - pkgworkflows "github.com/smartcontractkit/chainlink-common/pkg/workflows" "github.com/smartcontractkit/chainlink-evm/pkg/config" "github.com/smartcontractkit/chainlink/v2/core/services/job" "github.com/smartcontractkit/chainlink/v2/core/services/relay" @@ -313,109 +311,3 @@ func TestOCR2OracleSpec(t *testing.T) { }) }) } - -func TestWorkflowSpec_Validate(t *testing.T) { - if testing.Short() { - t.Skip("too slow for testing.Short") - } - - type fields struct { - Workflow string - } - tests := []struct { - name string - fields fields - wantWorkflowOwner string - wantWorkflowName string - - wantError bool - }{ - { - name: "valid", - fields: fields{ - Workflow: pkgworkflows.WFYamlSpec(t, "workflow01", "0x0123456789012345678901234567890123456789"), - }, - wantWorkflowOwner: "0123456789012345678901234567890123456789", // the workflow job spec strips the 0x prefix to limit to 40 characters - wantWorkflowName: "workflow01", - }, - { - name: "valid no name", - fields: fields{ - Workflow: pkgworkflows.WFYamlSpec(t, "", "0x0123456789012345678901234567890123456789"), - }, - wantWorkflowOwner: "0123456789012345678901234567890123456789", // the workflow job spec strips the 0x prefix to limit to 40 characters - wantWorkflowName: "", - }, - { - name: "valid no owner", - fields: fields{ - Workflow: pkgworkflows.WFYamlSpec(t, "workflow01", ""), - }, - wantWorkflowOwner: "", - wantWorkflowName: "workflow01", - }, - { - name: "invalid ", - fields: fields{ - Workflow: "garbage", - }, - wantError: true, - }, - } - for _, tt := range tests { - t.Run(tt.name, func(t *testing.T) { - w := &job.WorkflowSpec{ - Workflow: tt.fields.Workflow, - } - err := w.Validate(t.Context()) - require.Equal(t, tt.wantError, err != nil) - if !tt.wantError { - assert.NotEmpty(t, w.WorkflowID) - assert.Equal(t, tt.wantWorkflowOwner, w.WorkflowOwner) - assert.Equal(t, tt.wantWorkflowName, w.WorkflowName) - } - }) - } - - t.Run("WASM can validate", func(t *testing.T) { - configLocation := "testdata/config.json" - - w := &job.WorkflowSpec{ - Workflow: createTestBinary(t), - SpecType: job.WASMFile, - Config: configLocation, - } - - err := w.Validate(t.Context()) - require.NoError(t, err) - require.NotEmpty(t, w.WorkflowID) - }) - - t.Run("WASM can validate from TOML", func(t *testing.T) { - const wasmWorkflowTomlTemplate = ` - workflow_owner = "%s" - workflow_name = "%s" - spec_type = "%s" - workflow = "%s" - config = "%s" - ` - configLocation := "testdata/config.json" - tomlSpec := fmt.Sprintf( - wasmWorkflowTomlTemplate, - "0x0123456789012345678901234567890123456788", - "wf-2", - job.WASMFile, - createTestBinary(t), - configLocation, - ) - var w job.WorkflowSpec - err := toml.Unmarshal([]byte(tomlSpec), &w) - require.NoError(t, err) - - err = w.Validate(t.Context()) - require.NoError(t, err) - require.NotEmpty(t, w.WorkflowID) - assert.Equal(t, "0123456789012345678901234567890123456788", w.WorkflowOwner) - assert.Equal(t, "wf-2", w.WorkflowName) - }) -} diff --git a/core/services/job/orm.go b/core/services/job/orm.go index ed4d471715b..6f83d51d7a1 100644 --- a/core/services/job/orm.go +++ b/core/services/job/orm.go @@ -75,7 +75,6 @@ type ORM interface { DataSource() sqlutil.DataSource WithDataSource(source sqlutil.DataSource) ORM - FindJobIDByWorkflow(ctx context.Context, spec WorkflowSpec) (int32, error) // TODO rename function to indicate it is CCIP-specific, not generic? FindJobIDByCapabilityNameAndVersion(ctx context.Context, spec CCIPSpec) (int32, error) FindStandardCapabilityJobID(ctx context.Context, spec StandardCapabilitiesSpec) (int32, error) @@ -174,7 +173,7 @@ var ErrJobTypeRemoved = fmt.Errorf("job type has been removed and is no longer s func (o *orm) CreateJob(ctx context.Context, jb *Job) error { // Permanently removed job types: reject all new submissions regardless of // which code path reaches here (REST API, GraphQL, feeds manager, etc.). - if jb.Type == DirectRequest || jb.Type == FluxMonitor || jb.Type == Webhook { + if jb.Type == DirectRequest || jb.Type == FluxMonitor || jb.Type == Webhook || jb.Type == Workflow { return fmt.Errorf("cannot create job of type %q: %w", jb.Type, ErrJobTypeRemoved) } @@ -418,15 +417,6 @@ func (o *orm) CreateJob(ctx context.Context, jb *Job) error { jb.GatewaySpecID = &specID case Stream: // 'stream' type has no associated spec, nothing to do here - case Workflow: - sql := `INSERT INTO workflow_specs (workflow, workflow_id, workflow_owner, workflow_name, binary_url, config_url, secrets_id, created_at, updated_at, spec_type, config) - VALUES (:workflow, :workflow_id, :workflow_owner, :workflow_name, :binary_url, :config_url, :secrets_id, NOW(), NOW(), :spec_type, :config) - RETURNING id;` - specID, err := tx.prepareQuerySpecID(ctx, sql, jb.WorkflowSpec) - if err != nil { - return fmt.Errorf("failed to create WorkflowSpec for jobSpec given %v: %w", *jb.WorkflowSpec, err) - } - jb.WorkflowSpecID = &specID case StandardCapabilities: sql := `INSERT INTO standardcapabilities_specs (command, config, oracle_factory, created_at, updated_at) VALUES (:command, :config, :oracle_factory, NOW(), NOW()) @@ -752,7 +742,6 @@ func (o *orm) DeleteJob(ctx context.Context, id int32, jobType Type) error { Bootstrap: `DELETE FROM bootstrap_specs WHERE id IN (SELECT bootstrap_spec_id FROM deleted_jobs)`, BlockHeaderFeeder: `DELETE FROM block_header_feeder_specs WHERE id IN (SELECT block_header_feeder_spec_id FROM deleted_jobs)`, Gateway: `DELETE FROM gateway_specs WHERE id IN (SELECT gateway_spec_id FROM deleted_jobs)`, - Workflow: `DELETE FROM workflow_specs WHERE id in (SELECT workflow_spec_id FROM deleted_jobs)`, StandardCapabilities: `DELETE FROM standardcapabilities_specs WHERE id in (SELECT standard_capabilities_spec_id FROM deleted_jobs)`, CCIP: `DELETE FROM ccip_specs WHERE id in (SELECT ccip_spec_id FROM deleted_jobs)`, CCVCommitteeVerifier: `DELETE FROM ccv_committee_verifier_specs WHERE id IN (SELECT ccv_committee_verifier_spec_id FROM deleted_jobs)`, @@ -1139,23 +1128,6 @@ func (o *orm) FindJobIDsWithBridge(ctx context.Context, name string) (jids []int return jids, err } -func (o *orm) FindJobIDByWorkflow(ctx context.Context, spec WorkflowSpec) (jobID int32, err error) { - stmt := ` -SELECT jobs.id FROM jobs -INNER JOIN workflow_specs ws on jobs.workflow_spec_id = ws.id AND ws.workflow_owner = $1 AND ws.workflow_name = $2 -` - err = o.ds.GetContext(ctx, &jobID, stmt, spec.WorkflowOwner, spec.WorkflowName) - if err != nil { - if !errors.Is(err, sql.ErrNoRows) { - err = fmt.Errorf("error searching for job by workflow (owner,name) ('%s','%s'): %w", spec.WorkflowOwner, spec.WorkflowName, err) - } - err = fmt.Errorf("FindJobIDByWorkflow failed: %w", err) - return jobID, err - } - - return jobID, err -} - func (o *orm) FindJobIDByCapabilityNameAndVersion(ctx context.Context, spec CCIPSpec) (jobID int32, err error) { stmt := ` SELECT jobs.id FROM jobs diff --git a/core/services/job/testdata/config.json b/core/services/job/testdata/config.json deleted file mode 100644 index 57cba7a61c7..00000000000 --- a/core/services/job/testdata/config.json +++ /dev/null @@ -1,4 +0,0 @@ -{ - "Owner": "owner", - "Name": "name" -} \ No newline at end of file diff --git a/core/services/job/testdata/wasm/test_workflow_spec.go b/core/services/job/testdata/wasm/test_workflow_spec.go deleted file mode 100644 index 477ba097ca3..00000000000 --- a/core/services/job/testdata/wasm/test_workflow_spec.go +++ /dev/null @@ -1,25 +0,0 @@ -//go:build wasip1 - -package main - -import ( - "github.com/smartcontractkit/chainlink-common/pkg/workflows/wasm" - - "github.com/smartcontractkit/chainlink-common/pkg/capabilities/cli/cmd/testdata/fixtures/capabilities/basictrigger" - "github.com/smartcontractkit/chainlink-common/pkg/workflows/sdk" -) - -func BuildWorkflow(config []byte) *sdk.WorkflowSpecFactory { - workflow := sdk.NewWorkflowSpecFactory() - - triggerCfg := basictrigger.TriggerConfig{Name: "trigger", Number: 100} - _ = triggerCfg.New(workflow) - - return workflow -} - -func main() { - runner := wasm.NewRunner() - workflow := BuildWorkflow(runner.Config()) - runner.Run(workflow) -} diff --git a/core/services/job/wasm_file_spec_factory.go b/core/services/job/wasm_file_spec_factory.go deleted file mode 100644 index dbdd41f34be..00000000000 --- a/core/services/job/wasm_file_spec_factory.go +++ /dev/null @@ -1,114 +0,0 @@ -package job - -import ( - "bytes" - "context" - "crypto/sha256" - "encoding/hex" - "errors" - "fmt" - "io" - "os" - "path" - "strings" - - "github.com/andybalholm/brotli" - - "github.com/smartcontractkit/chainlink-common/pkg/workflows/sdk" - "github.com/smartcontractkit/chainlink-common/pkg/workflows/wasm/host" - "github.com/smartcontractkit/chainlink/v2/core/logger" -) - -type WasmFileSpecFactory struct{} - -func (w WasmFileSpecFactory) Spec(ctx context.Context, workflow, configLocation string) (sdk.WorkflowSpec, []byte, string, error) { - config, err := w.Config(ctx, configLocation) - if err != nil { - return sdk.WorkflowSpec{}, nil, "", err - } - - compressedBinary, sha, err := w.rawSpecAndSha(workflow, config) - if err != nil { - return sdk.WorkflowSpec{}, nil, "", err - } - - moduleConfig := &host.ModuleConfig{Logger: logger.NullLogger} - spec, err := host.GetWorkflowSpec(ctx, moduleConfig, compressedBinary, config) - if err != nil { - return sdk.WorkflowSpec{}, nil, "", err - } else if spec == nil { - return sdk.WorkflowSpec{}, nil, "", errors.New("workflow spec not found when running wasm") - } - - return *spec, compressedBinary, sha, nil -} - -func (w WasmFileSpecFactory) RawSpec(ctx context.Context, workflow, configLocation string) ([]byte, error) { - config, err := w.Config(ctx, configLocation) - if err != nil { - return nil, err - } - - raw, _, err := w.rawSpecAndSha(workflow, config) - return raw, err -} - -func (w WasmFileSpecFactory) Config(_ context.Context, configLocation string) ([]byte, error) { - config, err := os.ReadFile(configLocation) - if err != nil { - return nil, err - } - - return config, nil -} - -// rawSpecAndSha returns the brotli compressed version of the raw wasm file, alongside the sha256 hash of the raw wasm file -func (w WasmFileSpecFactory) rawSpecAndSha(wf string, config []byte) ([]byte, string, error) { - read, err := os.ReadFile(wf) - if err != nil { - return nil, "", err - } - - extension := strings.ToLower(path.Ext(wf)) - switch extension { - case ".wasm", "": - return w.rawSpecAndShaFromWasm(read, config) - case ".br": - return w.rawSpecAndShaFromBrotli(read, config) - default: - return nil, "", fmt.Errorf("unsupported file type %s", extension) - } -} - -func (w WasmFileSpecFactory) rawSpecAndShaFromBrotli(wasm, config []byte) ([]byte, string, error) { - brr := brotli.NewReader(bytes.NewReader(wasm)) - rawWasm, err := io.ReadAll(brr) - if err != nil { - return nil, "", err - } - - return wasm, w.sha(rawWasm, config), nil -} - -func (w WasmFileSpecFactory) rawSpecAndShaFromWasm(wasm, config []byte) ([]byte, string, error) { - var b bytes.Buffer - bwr := brotli.NewWriter(&b) - if _, err := bwr.Write(wasm); err != nil { - return nil, "", err - } - - if err := bwr.Close(); err != nil { - return nil, "", err - } - - return b.Bytes(), w.sha(wasm, config), nil -} - -func (w WasmFileSpecFactory) sha(wasm, config []byte) string { - sum := sha256.New() - sum.Write(wasm) - sum.Write(config) - return hex.EncodeToString(sum.Sum(nil)) -} - -var _ WorkflowSpecFactory = (*WasmFileSpecFactory)(nil) diff --git a/core/services/job/wasm_file_spec_factory_test.go b/core/services/job/wasm_file_spec_factory_test.go deleted file mode 100644 index 9735c830ecb..00000000000 --- a/core/services/job/wasm_file_spec_factory_test.go +++ /dev/null @@ -1,102 +0,0 @@ -package job_test - -import ( - "bytes" - "crypto/sha256" - "encoding/hex" - "os" - "os/exec" - "strings" - "testing" - - "github.com/andybalholm/brotli" - "github.com/stretchr/testify/assert" - "github.com/stretchr/testify/require" - - "github.com/smartcontractkit/chainlink-common/pkg/workflows/wasm/host" - "github.com/smartcontractkit/chainlink/v2/core/logger" - "github.com/smartcontractkit/chainlink/v2/core/services/job" -) - -func TestWasmFileSpecFactory(t *testing.T) { - if testing.Short() { - t.Skip("too slow for testing.Short") - } - - binaryLocation := createTestBinary(t) - configLocation := "testdata/config.json" - config, err := os.ReadFile(configLocation) - require.NoError(t, err) - - rawBinary, err := os.ReadFile(binaryLocation) - require.NoError(t, err) - - b := bytes.Buffer{} - bwr := brotli.NewWriter(&b) - _, err = bwr.Write(rawBinary) - require.NoError(t, err) - - require.NoError(t, bwr.Close()) - - t.Run("Raw binary", func(t *testing.T) { - ctx := t.Context() - factory := job.WasmFileSpecFactory{} - actual, rawSpec, actualSha, err2 := factory.Spec(t.Context(), binaryLocation, configLocation) - require.NoError(t, err2) - - expected, err2 := host.GetWorkflowSpec(ctx, &host.ModuleConfig{Logger: logger.NullLogger, IsUncompressed: true}, rawBinary, config) - require.NoError(t, err2) - - expectedSha := sha256.New() - expectedSha.Write(rawBinary) - expectedSha.Write(config) - require.Equal(t, hex.EncodeToString(expectedSha.Sum(nil)), actualSha) - - require.Equal(t, *expected, actual) - - assert.Equal(t, b.Bytes(), rawSpec) - }) - - t.Run("Compressed binary", func(t *testing.T) { - ctx := t.Context() - brLoc := strings.Replace(binaryLocation, ".wasm", ".br", 1) - compressedBytes := b.Bytes() - require.NoError(t, os.WriteFile(brLoc, compressedBytes, 0o600)) - - factory := job.WasmFileSpecFactory{} - actual, rawSpec, actualSha, err2 := factory.Spec(t.Context(), brLoc, configLocation) - require.NoError(t, err2) - - expected, err2 := host.GetWorkflowSpec(ctx, &host.ModuleConfig{Logger: logger.NullLogger, IsUncompressed: true}, rawBinary, config) - require.NoError(t, err2) - - expectedSha := sha256.New() - expectedSha.Write(rawBinary) - expectedSha.Write(config) - require.Equal(t, hex.EncodeToString(expectedSha.Sum(nil)), actualSha) - - require.Equal(t, *expected, actual) - - assert.Equal(t, b.Bytes(), rawSpec) - }) - - t.Run("Config", func(t *testing.T) { - factory := job.WasmFileSpecFactory{} - actual, err3 := factory.Config(t.Context(), configLocation) - require.NoError(t, err3) - - assert.Equal(t, config, actual) - }) -} - -func createTestBinary(t *testing.T) string { - const testBinaryLocation = "testdata/wasm/testmodule.wasm" - - cmd := exec.CommandContext(t.Context(), "go", "build", "-o", testBinaryLocation, "github.com/smartcontractkit/chainlink/v2/core/services/job/testdata/wasm") - cmd.Env = append(os.Environ(), "GOOS=wasip1", "GOARCH=wasm") - - output, err := cmd.CombinedOutput() - require.NoError(t, err, string(output)) - - return testBinaryLocation -} diff --git a/core/services/job/workflow_spec_factory.go b/core/services/job/workflow_spec_factory.go deleted file mode 100644 index c799b69823e..00000000000 --- a/core/services/job/workflow_spec_factory.go +++ /dev/null @@ -1,19 +0,0 @@ -package job - -import ( - "context" - - "github.com/smartcontractkit/chainlink-common/pkg/workflows/sdk" -) - -type WorkflowSpecFactory interface { - Spec(ctx context.Context, workflow, config string) (sdk.WorkflowSpec, []byte, string, error) - RawSpec(ctx context.Context, workflow, config string) ([]byte, error) - Config(ctx context.Context, config string) ([]byte, error) -} - -var workflowSpecFactories = map[WorkflowSpecType]WorkflowSpecFactory{ - YamlSpec: YAMLSpecFactory{}, - WASMFile: WasmFileSpecFactory{}, - DefaultSpecType: YAMLSpecFactory{}, -} diff --git a/core/services/job/yaml_spec_factory.go b/core/services/job/yaml_spec_factory.go deleted file mode 100644 index 92df8a6ed0c..00000000000 --- a/core/services/job/yaml_spec_factory.go +++ /dev/null @@ -1,27 +0,0 @@ -package job - -import ( - "context" - "crypto/sha256" - "fmt" - - "github.com/smartcontractkit/chainlink-common/pkg/workflows" - "github.com/smartcontractkit/chainlink-common/pkg/workflows/sdk" -) - -type YAMLSpecFactory struct{} - -var _ WorkflowSpecFactory = (*YAMLSpecFactory)(nil) - -func (y YAMLSpecFactory) Spec(_ context.Context, workflow, _ string) (sdk.WorkflowSpec, []byte, string, error) { - spec, err := workflows.ParseWorkflowSpecYaml(workflow) - return spec, []byte(workflow), fmt.Sprintf("%x", sha256.Sum256([]byte(workflow))), err -} - -func (y YAMLSpecFactory) RawSpec(_ context.Context, workflow, _ string) ([]byte, error) { - return []byte(workflow), nil -} - -func (y YAMLSpecFactory) Config(_ context.Context, config string) ([]byte, error) { - return []byte(config), nil -} diff --git a/core/services/job/yaml_spec_factory_test.go b/core/services/job/yaml_spec_factory_test.go deleted file mode 100644 index d6752be3d37..00000000000 --- a/core/services/job/yaml_spec_factory_test.go +++ /dev/null @@ -1,86 +0,0 @@ -package job_test - -import ( - "crypto/sha256" - "fmt" - "testing" - - "github.com/stretchr/testify/assert" - "github.com/stretchr/testify/require" - - commonworkflows "github.com/smartcontractkit/chainlink-common/pkg/workflows" - "github.com/smartcontractkit/chainlink/v2/core/services/job" -) - -const anyYamlSpec = ` -name: "wf-name" -owner: "0x00000000000000000000000000000000000000aa" -triggers: - - id: "mercury-trigger@1.0.0" - config: - feedIds: - - "0x1111111111111111111100000000000000000000000000000000000000000000" - - "0x2222222222222222222200000000000000000000000000000000000000000000" - - "0x3333333333333333333300000000000000000000000000000000000000000000" - -consensus: - - id: "offchain_reporting@2.0.0" - ref: "evm_median" - inputs: - observations: - - "$(trigger.outputs)" - config: - aggregation_method: "data_feeds_2_0" - aggregation_config: - "0x1111111111111111111100000000000000000000000000000000000000000000": - deviation: "0.001" - heartbeat: 3600 - "0x2222222222222222222200000000000000000000000000000000000000000000": - deviation: "0.001" - heartbeat: 3600 - "0x3333333333333333333300000000000000000000000000000000000000000000": - deviation: "0.001" - heartbeat: 3600 - encoder: "EVM" - encoder_config: - abi: "mercury_reports bytes[]" - -targets: - - id: "write_polygon-testnet-mumbai@3.0.0" - inputs: - report: "$(evm_median.outputs.report)" - config: - address: "0x3F3554832c636721F1fD1822Ccca0354576741Ef" - params: ["$(report)"] - abi: "receive(report bytes)" - - id: "write_ethereum-testnet-sepolia@4.0.0" - inputs: - report: "$(evm_median.outputs.report)" - config: - address: "0x54e220867af6683aE6DcBF535B4f952cB5116510" - params: ["$(report)"] - abi: "receive(report bytes)" -` - -func TestYamlSpecFactory_GetSpec(t *testing.T) { - t.Parallel() - - actual, raw, actualSha, err := job.YAMLSpecFactory{}.Spec(t.Context(), anyYamlSpec, "") - require.NoError(t, err) - - expected, err := commonworkflows.ParseWorkflowSpecYaml(anyYamlSpec) - require.NoError(t, err) - - require.Equal(t, expected, actual) - assert.Equal(t, fmt.Sprintf("%x", sha256.Sum256([]byte(anyYamlSpec))), actualSha) - assert.YAMLEq(t, anyYamlSpec, string(raw)) -} - -func TestYamlSpecFactory_Config(t *testing.T) { - t.Parallel() - - config := "config" - actual, err := job.YAMLSpecFactory{}.Config(t.Context(), config) - require.NoError(t, err) - assert.Equal(t, []byte(config), actual) -} diff --git a/core/services/workflows/artifacts/orm.go b/core/services/workflows/artifacts/orm.go deleted file mode 100644 index b8dde5a9fea..00000000000 --- a/core/services/workflows/artifacts/orm.go +++ /dev/null @@ -1,437 +0,0 @@ -package artifacts - -import ( - "context" - "database/sql" - "errors" - "fmt" - "time" - - "github.com/jmoiron/sqlx" - - "github.com/smartcontractkit/chainlink-common/pkg/logger" - "github.com/smartcontractkit/chainlink-common/pkg/sqlutil" - "github.com/smartcontractkit/chainlink/v2/core/services/job" - "github.com/smartcontractkit/chainlink/v2/core/utils/crypto" -) - -type WorkflowSecretsDS interface { - // GetSecretsURLByID returns the secrets URL for the given ID. - GetSecretsURLByID(ctx context.Context, id int64) (string, error) - - // GetSecretsURLByID returns the secrets URL for the given ID. - GetSecretsURLByHash(ctx context.Context, hash string) (string, error) - - // GetContents returns the contents of the secret at the given plain URL. - GetContents(ctx context.Context, url string) (string, error) - - // GetContentsByHash returns the contents of the secret at the given hashed URL. - GetContentsByHash(ctx context.Context, hash string) (string, error) - - // GetContentsByWorkflowID returns the contents and secrets_url of the secret for the given workflow. - GetContentsByWorkflowID(ctx context.Context, workflowID string) (string, string, error) - - // GetSecretsURLHash returns the keccak256 hash of the owner and secrets URL. - GetSecretsURLHash(owner, secretsURL []byte) ([]byte, error) - - // Update updates the contents of the secrets at the given plain URL or inserts a new record if not found. - Update(ctx context.Context, secretsURL, contents string) (int64, error) - - Create(ctx context.Context, secretsURL, hash, contents string) (int64, error) -} - -type WorkflowSpecsDS interface { - // UpsertWorkflowSpec inserts or updates a workflow spec. Updates on conflict of workflow name - // and owner - UpsertWorkflowSpec(ctx context.Context, spec *job.WorkflowSpec) (int64, error) - - // UpsertWorkflowSpecWithSecrets inserts or updates a workflow spec with secrets in a transaction. - // Updates on conflict of workflow name and owner. - UpsertWorkflowSpecWithSecrets(ctx context.Context, spec *job.WorkflowSpec, url, hash, contents string) (int64, error) - - // GetWorkflowSpec returns the workflow spec for the given owner and name. - GetWorkflowSpec(ctx context.Context, owner, name string) (*job.WorkflowSpec, error) - - // DeleteWorkflowSpec deletes the workflow spec for the given owner and name. - DeleteWorkflowSpec(ctx context.Context, owner, name string) error - - // GetWorkflowSpecByID returns the workflow spec for the given workflowID. - GetWorkflowSpecByID(ctx context.Context, id string) (*job.WorkflowSpec, error) -} - -type ORM interface { - WorkflowSecretsDS - WorkflowSpecsDS -} - -type WorkflowRegistryDS = ORM - -type orm struct { - ds sqlutil.DataSource - lggr logger.Logger -} - -var _ WorkflowRegistryDS = (*orm)(nil) - -func NewWorkflowRegistryDS(ds sqlutil.DataSource, lggr logger.Logger) *orm { - return &orm{ - ds: ds, - lggr: lggr, - } -} - -func (orm *orm) GetSecretsURLByID(ctx context.Context, id int64) (string, error) { - var secretsURL string - err := orm.ds.GetContext(ctx, &secretsURL, - `SELECT secrets_url FROM workflow_secrets WHERE workflow_secrets.id = $1`, - id, - ) - - return secretsURL, err -} - -func (orm *orm) GetSecretsURLByHash(ctx context.Context, hash string) (string, error) { - var secretsURL string - err := orm.ds.GetContext(ctx, &secretsURL, - `SELECT secrets_url FROM workflow_secrets WHERE workflow_secrets.secrets_url_hash = $1`, - hash, - ) - - return secretsURL, err -} - -func (orm *orm) GetContentsByHash(ctx context.Context, hash string) (string, error) { - var contents string - err := orm.ds.GetContext(ctx, &contents, - `SELECT contents - FROM workflow_secrets - WHERE secrets_url_hash = $1`, - hash, - ) - - if err != nil { - return "", err // Return an empty Artifact struct and the error - } - - return contents, nil // Return the populated Artifact struct -} - -func (orm *orm) GetContents(ctx context.Context, url string) (string, error) { - var contents string - err := orm.ds.GetContext(ctx, &contents, - `SELECT contents - FROM workflow_secrets - WHERE secrets_url = $1`, - url, - ) - - if err != nil { - return "", err // Return an empty Artifact struct and the error - } - - return contents, nil // Return the populated Artifact struct -} - -type Int struct { - sql.NullInt64 -} - -type joinRecord struct { - SecretsID sql.NullString `db:"wspec_secrets_id"` - SecretsURLHash sql.NullString `db:"wsec_secrets_url_hash"` - Contents sql.NullString `db:"wsec_contents"` -} - -var ErrEmptySecrets = errors.New("secrets field is empty") - -// GetContentsByWorkflowID joins the workflow_secrets on the workflow_specs table and gets -// the associated secrets contents. -func (orm *orm) GetContentsByWorkflowID(ctx context.Context, workflowID string) (string, string, error) { - var jr joinRecord - err := orm.ds.GetContext( - ctx, - &jr, - `SELECT wsec.secrets_url_hash AS wsec_secrets_url_hash, wsec.contents AS wsec_contents, wspec.secrets_id AS wspec_secrets_id - FROM workflow_specs AS wspec - LEFT JOIN - workflow_secrets AS wsec ON wspec.secrets_id = wsec.id - WHERE wspec.workflow_id = $1`, - workflowID, - ) - if err != nil { - return "", "", err - } - - if !jr.SecretsID.Valid { - return "", "", ErrEmptySecrets - } - - if jr.Contents.String == "" { - return "", "", ErrEmptySecrets - } - - return jr.SecretsURLHash.String, jr.Contents.String, nil -} - -// Update updates the secrets content at the given hash or inserts a new record if not found. -func (orm *orm) Update(ctx context.Context, hash, contents string) (int64, error) { - var id int64 - err := orm.ds.QueryRowxContext(ctx, - `INSERT INTO workflow_secrets (secrets_url_hash, contents) - VALUES ($1, $2) - ON CONFLICT (secrets_url_hash) DO UPDATE - SET secrets_url_hash = EXCLUDED.secrets_url_hash, contents = EXCLUDED.contents - RETURNING id`, - hash, contents, - ).Scan(&id) - - if err != nil { - return 0, err - } - - return id, nil -} - -// Update updates the secrets content at the given hash or inserts a new record if not found. -func (orm *orm) Create(ctx context.Context, url, hash, contents string) (int64, error) { - var id int64 - err := orm.ds.QueryRowxContext(ctx, - `INSERT INTO workflow_secrets (secrets_url, secrets_url_hash, contents) - VALUES ($1, $2, $3) - RETURNING id`, - url, hash, contents, - ).Scan(&id) - - if err != nil { - return 0, err - } - - return id, nil -} - -func (orm *orm) GetSecretsURLHash(owner, secretsURL []byte) ([]byte, error) { - return crypto.Keccak256(append(owner, secretsURL...)) -} - -func (orm *orm) UpsertWorkflowSpec(ctx context.Context, spec *job.WorkflowSpec) (int64, error) { - var id int64 - err := sqlutil.TransactDataSource(ctx, orm.ds, nil, func(tx sqlutil.DataSource) error { - _, txErr := tx.ExecContext( - ctx, - `DELETE FROM workflow_specs WHERE workflow_owner = $1 AND workflow_name = $2 AND workflow_id != $3`, - spec.WorkflowOwner, - spec.WorkflowName, - spec.WorkflowID, - ) - if txErr != nil { - return fmt.Errorf("failed to clean up previous workflow specs: %w", txErr) - } - - query := ` - INSERT INTO workflow_specs ( - workflow, - config, - workflow_id, - workflow_owner, - workflow_name, - status, - binary_url, - config_url, - secrets_id, - created_at, - updated_at, - spec_type - ) VALUES ( - :workflow, - :config, - :workflow_id, - :workflow_owner, - :workflow_name, - :status, - :binary_url, - :config_url, - :secrets_id, - :created_at, - :updated_at, - :spec_type - ) ON CONFLICT (workflow_owner, workflow_name) DO UPDATE - SET - workflow = EXCLUDED.workflow, - config = EXCLUDED.config, - workflow_id = EXCLUDED.workflow_id, - workflow_owner = EXCLUDED.workflow_owner, - workflow_name = EXCLUDED.workflow_name, - status = EXCLUDED.status, - binary_url = EXCLUDED.binary_url, - config_url = EXCLUDED.config_url, - secrets_id = EXCLUDED.secrets_id, - created_at = EXCLUDED.created_at, - updated_at = EXCLUDED.updated_at, - spec_type = EXCLUDED.spec_type - RETURNING id - ` - - now := time.Now().UTC() - spec.UpdatedAt = now - if spec.CreatedAt.IsZero() { - spec.CreatedAt = now - } - q, args, namedErr := sqlx.Named(query, spec) - if namedErr != nil { - return namedErr - } - q = sqlx.Rebind(sqlx.DOLLAR, q) - return tx.QueryRowxContext(ctx, q, args...).Scan(&id) - }) - - return id, err -} - -func (orm *orm) UpsertWorkflowSpecWithSecrets( - ctx context.Context, - spec *job.WorkflowSpec, url, hash, contents string) (int64, error) { - var id int64 - err := sqlutil.TransactDataSource(ctx, orm.ds, nil, func(tx sqlutil.DataSource) error { - var sid int64 - txErr := tx.QueryRowxContext(ctx, - `INSERT INTO workflow_secrets (secrets_url, secrets_url_hash, contents) - VALUES ($1, $2, $3) - ON CONFLICT (secrets_url_hash) DO UPDATE - SET - secrets_url_hash = EXCLUDED.secrets_url_hash, - contents = EXCLUDED.contents, - secrets_url = EXCLUDED.secrets_url - RETURNING id`, - url, hash, contents, - ).Scan(&sid) - - if txErr != nil { - return fmt.Errorf("failed to create workflow secrets: %w", txErr) - } - - _, queryErr := tx.ExecContext( - ctx, - `DELETE FROM workflow_specs WHERE workflow_owner = $1 AND workflow_name = $2 AND workflow_id != $3`, - spec.WorkflowOwner, - spec.WorkflowName, - spec.WorkflowID, - ) - if queryErr != nil { - return fmt.Errorf("failed to clean up previous workflow specs: %w", queryErr) - } - - spec.SecretsID = sql.NullInt64{Int64: sid, Valid: true} - - query := ` - INSERT INTO workflow_specs ( - workflow, - config, - workflow_id, - workflow_owner, - workflow_name, - status, - binary_url, - config_url, - secrets_id, - created_at, - updated_at, - spec_type - ) VALUES ( - :workflow, - :config, - :workflow_id, - :workflow_owner, - :workflow_name, - :status, - :binary_url, - :config_url, - :secrets_id, - :created_at, - :updated_at, - :spec_type - ) ON CONFLICT (workflow_owner, workflow_name) DO UPDATE - SET - workflow = EXCLUDED.workflow, - config = EXCLUDED.config, - workflow_id = EXCLUDED.workflow_id, - workflow_owner = EXCLUDED.workflow_owner, - workflow_name = EXCLUDED.workflow_name, - status = EXCLUDED.status, - binary_url = EXCLUDED.binary_url, - config_url = EXCLUDED.config_url, - created_at = EXCLUDED.created_at, - updated_at = EXCLUDED.updated_at, - spec_type = EXCLUDED.spec_type, - secrets_id = EXCLUDED.secrets_id - RETURNING id - ` - - now := time.Now().UTC() - spec.UpdatedAt = now - if spec.CreatedAt.IsZero() { - spec.CreatedAt = now - } - q, args, namedErr := sqlx.Named(query, spec) - if namedErr != nil { - return namedErr - } - q = sqlx.Rebind(sqlx.DOLLAR, q) - return tx.QueryRowxContext(ctx, q, args...).Scan(&id) - }) - return id, err -} - -func (orm *orm) GetWorkflowSpec(ctx context.Context, owner, name string) (*job.WorkflowSpec, error) { - query := ` - SELECT * - FROM workflow_specs - WHERE workflow_owner = $1 AND workflow_name = $2 - ` - - var spec job.WorkflowSpec - err := orm.ds.GetContext(ctx, &spec, query, owner, name) - if err != nil { - return nil, err - } - - return &spec, nil -} - -func (orm *orm) GetWorkflowSpecByID(ctx context.Context, id string) (*job.WorkflowSpec, error) { - query := ` - SELECT * - FROM workflow_specs - WHERE workflow_id = $1 - ` - - var spec job.WorkflowSpec - err := orm.ds.GetContext(ctx, &spec, query, id) - if err != nil { - return nil, err - } - - return &spec, nil -} - -func (orm *orm) DeleteWorkflowSpec(ctx context.Context, owner, name string) error { - query := ` - DELETE FROM workflow_specs - WHERE workflow_owner = $1 AND workflow_name = $2 - ` - - result, err := orm.ds.ExecContext(ctx, query, owner, name) - if err != nil { - return err - } - - rowsAffected, err := result.RowsAffected() - if err != nil { - return err - } - - if rowsAffected == 0 { - return sql.ErrNoRows // No spec deleted - } - - return nil -} diff --git a/core/services/workflows/artifacts/orm_test.go b/core/services/workflows/artifacts/orm_test.go deleted file mode 100644 index 0c04af66441..00000000000 --- a/core/services/workflows/artifacts/orm_test.go +++ /dev/null @@ -1,489 +0,0 @@ -package artifacts - -import ( - "database/sql" - "encoding/hex" - "testing" - "time" - - "github.com/google/uuid" - "github.com/stretchr/testify/assert" - "github.com/stretchr/testify/require" - - "github.com/smartcontractkit/chainlink/v2/core/internal/testutils/pgtest" - "github.com/smartcontractkit/chainlink/v2/core/logger" - "github.com/smartcontractkit/chainlink/v2/core/services/job" - "github.com/smartcontractkit/chainlink/v2/core/utils/crypto" -) - -func TestWorkflowArtifactsORM_GetAndUpdate(t *testing.T) { - t.Parallel() - db := pgtest.NewSqlxDB(t) - ctx := t.Context() - lggr := logger.TestLogger(t) - orm := &orm{ds: db, lggr: lggr} - - giveURL := "https://example.com/" + uuid.New().String()[:8] - giveBytes, err := crypto.Keccak256([]byte(giveURL)) - require.NoError(t, err) - giveHash := hex.EncodeToString(giveBytes) - giveContent := "some contents" - - gotID, err := orm.Create(ctx, giveURL, giveHash, giveContent) - require.NoError(t, err) - - url, err := orm.GetSecretsURLByID(ctx, gotID) - require.NoError(t, err) - assert.Equal(t, giveURL, url) - - contents, err := orm.GetContents(ctx, giveURL) - require.NoError(t, err) - assert.Equal(t, "some contents", contents) - - contents, err = orm.GetContentsByHash(ctx, giveHash) - require.NoError(t, err) - assert.Equal(t, "some contents", contents) - - _, err = orm.Update(ctx, giveHash, "new contents") - require.NoError(t, err) - - contents, err = orm.GetContents(ctx, giveURL) - require.NoError(t, err) - assert.Equal(t, "new contents", contents) - - contents, err = orm.GetContentsByHash(ctx, giveHash) - require.NoError(t, err) - assert.Equal(t, "new contents", contents) -} - -func Test_UpsertWorkflowSpec(t *testing.T) { - t.Parallel() - db := pgtest.NewSqlxDB(t) - ctx := t.Context() - lggr := logger.TestLogger(t) - orm := &orm{ds: db, lggr: lggr} - - owner := "owner-" + uuid.New().String()[:8] - name := "name-" + uuid.New().String()[:8] - cid := "cid-" + uuid.New().String()[:8] - - t.Run("inserts new spec", func(t *testing.T) { //nolint:paralleltest // subtests share database setup - spec := &job.WorkflowSpec{ - Workflow: "test_workflow", - Config: "test_config", - WorkflowID: cid, - WorkflowOwner: owner, - WorkflowName: name, - Status: job.WorkflowSpecStatusActive, - BinaryURL: "http://example.com/binary/" + cid, - ConfigURL: "http://example.com/config/" + cid, - CreatedAt: time.Now(), - SpecType: job.WASMFile, - } - - _, err := orm.UpsertWorkflowSpec(ctx, spec) - require.NoError(t, err) - - // Verify the record exists in the database - var dbSpec job.WorkflowSpec - err = db.Get(&dbSpec, `SELECT * FROM workflow_specs WHERE workflow_owner = $1 AND workflow_name = $2`, spec.WorkflowOwner, spec.WorkflowName) - require.NoError(t, err) - require.Equal(t, spec.Workflow, dbSpec.Workflow) - }) - - t.Run("updates existing spec", func(t *testing.T) { //nolint:paralleltest // subtests share database setup - spec := &job.WorkflowSpec{ - Workflow: "test_workflow", - Config: "test_config", - WorkflowID: cid, - WorkflowOwner: owner, - WorkflowName: name, - Status: job.WorkflowSpecStatusActive, - BinaryURL: "http://example.com/binary/" + cid, - ConfigURL: "http://example.com/config/" + cid, - CreatedAt: time.Now(), - SpecType: job.WASMFile, - } - - _, err := orm.UpsertWorkflowSpec(ctx, spec) - require.NoError(t, err) - - // Update the status - spec.Status = job.WorkflowSpecStatusPaused - - _, err = orm.UpsertWorkflowSpec(ctx, spec) - require.NoError(t, err) - - // Verify the record is updated in the database - var dbSpec job.WorkflowSpec - err = db.Get(&dbSpec, `SELECT * FROM workflow_specs WHERE workflow_owner = $1 AND workflow_name = $2`, spec.WorkflowOwner, spec.WorkflowName) - require.NoError(t, err) - require.Equal(t, spec.Config, dbSpec.Config) - require.Equal(t, spec.Status, dbSpec.Status) - }) -} - -func Test_DeleteWorkflowSpec(t *testing.T) { - t.Parallel() - db := pgtest.NewSqlxDB(t) - ctx := t.Context() - lggr := logger.TestLogger(t) - orm := &orm{ds: db, lggr: lggr} - - owner := "owner-" + uuid.New().String()[:8] - name := "name-" + uuid.New().String()[:8] - cid := "cid-" + uuid.New().String()[:8] - - t.Run("deletes a workflow spec", func(t *testing.T) { //nolint:paralleltest // subtests share database setup - spec := &job.WorkflowSpec{ - Workflow: "test_workflow", - Config: "test_config", - WorkflowID: cid, - WorkflowOwner: owner, - WorkflowName: name, - Status: job.WorkflowSpecStatusActive, - BinaryURL: "http://example.com/binary/" + cid, - ConfigURL: "http://example.com/config/" + cid, - CreatedAt: time.Now(), - SpecType: job.WASMFile, - } - - id, err := orm.UpsertWorkflowSpec(ctx, spec) - require.NoError(t, err) - require.NotZero(t, id) - - err = orm.DeleteWorkflowSpec(ctx, spec.WorkflowOwner, spec.WorkflowName) - require.NoError(t, err) - - // Verify the record is deleted from the database - var dbSpec job.WorkflowSpec - err = db.Get(&dbSpec, `SELECT * FROM workflow_specs WHERE id = $1`, id) - require.Error(t, err) - require.Equal(t, sql.ErrNoRows, err) - }) - - t.Run("fails if no workflow spec exists", func(t *testing.T) { //nolint:paralleltest // subtests share database setup - err := orm.DeleteWorkflowSpec(ctx, owner, name) - require.Error(t, err) - require.Equal(t, sql.ErrNoRows, err) - }) -} - -func Test_GetWorkflowSpec(t *testing.T) { - t.Parallel() - db := pgtest.NewSqlxDB(t) - ctx := t.Context() - lggr := logger.TestLogger(t) - orm := &orm{ds: db, lggr: lggr} - - owner := "owner-" + uuid.New().String()[:8] - name := "name-" + uuid.New().String()[:8] - cid := "cid-" + uuid.New().String()[:8] - - t.Run("gets a workflow spec", func(t *testing.T) { //nolint:paralleltest // subtests share database setup - spec := &job.WorkflowSpec{ - Workflow: "test_workflow", - Config: "test_config", - WorkflowID: cid, - WorkflowOwner: owner, - WorkflowName: name, - Status: job.WorkflowSpecStatusActive, - BinaryURL: "http://example.com/binary/" + cid, - ConfigURL: "http://example.com/config/" + cid, - CreatedAt: time.Now(), - SpecType: job.WASMFile, - } - - id, err := orm.UpsertWorkflowSpec(ctx, spec) - require.NoError(t, err) - require.NotZero(t, id) - - dbSpec, err := orm.GetWorkflowSpec(ctx, spec.WorkflowOwner, spec.WorkflowName) - require.NoError(t, err) - require.Equal(t, spec.Workflow, dbSpec.Workflow) - - err = orm.DeleteWorkflowSpec(ctx, spec.WorkflowOwner, spec.WorkflowName) - require.NoError(t, err) - }) - - t.Run("fails if no workflow spec exists", func(t *testing.T) { //nolint:paralleltest // subtests share database setup - dbSpec, err := orm.GetWorkflowSpec(ctx, owner, name) - require.Error(t, err) - require.Nil(t, dbSpec) - }) -} - -func Test_GetWorkflowSpecByID(t *testing.T) { - t.Parallel() - db := pgtest.NewSqlxDB(t) - ctx := t.Context() - lggr := logger.TestLogger(t) - orm := &orm{ds: db, lggr: lggr} - - owner := "owner-" + uuid.New().String()[:8] - name := "name-" + uuid.New().String()[:8] - cid := "cid-" + uuid.New().String()[:8] - - t.Run("gets a workflow spec by ID", func(t *testing.T) { //nolint:paralleltest // subtests share database setup - spec := &job.WorkflowSpec{ - Workflow: "test_workflow", - Config: "test_config", - WorkflowID: cid, - WorkflowOwner: owner, - WorkflowName: name, - Status: job.WorkflowSpecStatusActive, - BinaryURL: "http://example.com/binary/" + cid, - ConfigURL: "http://example.com/config/" + cid, - CreatedAt: time.Now(), - SpecType: job.WASMFile, - } - - id, err := orm.UpsertWorkflowSpec(ctx, spec) - require.NoError(t, err) - require.NotZero(t, id) - - dbSpec, err := orm.GetWorkflowSpecByID(ctx, spec.WorkflowID) - require.NoError(t, err) - require.Equal(t, spec.Workflow, dbSpec.Workflow) - - err = orm.DeleteWorkflowSpec(ctx, spec.WorkflowOwner, spec.WorkflowName) - require.NoError(t, err) - }) - - t.Run("fails if no workflow spec exists", func(t *testing.T) { //nolint:paralleltest // subtests share database setup - dbSpec, err := orm.GetWorkflowSpecByID(ctx, cid) - require.Error(t, err) - require.Nil(t, dbSpec) - }) -} - -func Test_GetContentsByWorkflowID(t *testing.T) { - t.Parallel() - db := pgtest.NewSqlxDB(t) - ctx := t.Context() - lggr := logger.TestLogger(t) - orm := &orm{ds: db, lggr: lggr} - - workflowID := "wf-" + uuid.New().String()[:8] - workflowOwner := "owner-" + uuid.New().String()[:8] - workflowName := "name-" + uuid.New().String()[:8] - - // workflow_id is missing - _, _, err := orm.GetContentsByWorkflowID(ctx, "doesnt-exist") - require.ErrorContains(t, err, "no rows in result set") - - // secrets_id is nil; should return EmptySecrets - _, err = orm.UpsertWorkflowSpec(ctx, &job.WorkflowSpec{ - Workflow: "", - Config: "", - WorkflowID: workflowID, - WorkflowOwner: workflowOwner, - WorkflowName: workflowName, - BinaryURL: "", - ConfigURL: "", - CreatedAt: time.Now(), - SpecType: job.DefaultSpecType, - }) - require.NoError(t, err) - - _, _, err = orm.GetContentsByWorkflowID(ctx, workflowID) - require.ErrorIs(t, err, ErrEmptySecrets) - - // retrieves the artifact if provided - giveURL := "https://example.com/" + uuid.New().String()[:8] - giveBytes, err := crypto.Keccak256([]byte(giveURL)) - require.NoError(t, err) - giveHash := hex.EncodeToString(giveBytes) - giveContent := "some contents" - - secretsID, err := orm.Create(ctx, giveURL, giveHash, giveContent) - require.NoError(t, err) - - // Use a new workflow ID to avoid conflicts - workflowID2 := "wf-" + uuid.New().String()[:8] - _, err = orm.UpsertWorkflowSpec(ctx, &job.WorkflowSpec{ - Workflow: "", - Config: "", - SecretsID: sql.NullInt64{Int64: secretsID, Valid: true}, - WorkflowID: workflowID2, - WorkflowOwner: workflowOwner, - WorkflowName: workflowName + "-2", - BinaryURL: "", - ConfigURL: "", - CreatedAt: time.Now(), - SpecType: job.DefaultSpecType, - }) - require.NoError(t, err) - _, err = orm.GetWorkflowSpec(ctx, workflowOwner, workflowName+"-2") - require.NoError(t, err) - - gotHash, gotContent, err := orm.GetContentsByWorkflowID(ctx, workflowID2) - require.NoError(t, err) - assert.Equal(t, giveHash, gotHash) - assert.Equal(t, giveContent, gotContent) -} - -func Test_GetContentsByWorkflowID_SecretsProvidedButEmpty(t *testing.T) { - t.Parallel() - db := pgtest.NewSqlxDB(t) - ctx := t.Context() - lggr := logger.TestLogger(t) - orm := &orm{ds: db, lggr: lggr} - - // workflow_id is missing - _, _, err := orm.GetContentsByWorkflowID(ctx, "doesnt-exist") - require.ErrorContains(t, err, "no rows in result set") - - // secrets_id is nil; should return EmptySecrets - workflowID := "wf-" + uuid.New().String()[:8] - workflowOwner := "owner-" + uuid.New().String()[:8] - workflowName := "name-" + uuid.New().String()[:8] - giveURL := "https://example.com/" + uuid.New().String()[:8] - giveBytes, err := crypto.Keccak256([]byte(giveURL)) - require.NoError(t, err) - giveHash := hex.EncodeToString(giveBytes) - giveContent := "" - _, err = orm.UpsertWorkflowSpecWithSecrets(ctx, &job.WorkflowSpec{ - Workflow: "", - Config: "", - WorkflowID: workflowID, - WorkflowOwner: workflowOwner, - WorkflowName: workflowName, - BinaryURL: "", - ConfigURL: "", - CreatedAt: time.Now(), - SpecType: job.DefaultSpecType, - }, giveURL, giveHash, giveContent) - require.NoError(t, err) - - _, _, err = orm.GetContentsByWorkflowID(ctx, workflowID) - require.ErrorIs(t, err, ErrEmptySecrets) -} - -func Test_UpsertWorkflowSpecWithSecrets(t *testing.T) { - t.Parallel() - db := pgtest.NewSqlxDB(t) - ctx := t.Context() - lggr := logger.TestLogger(t) - orm := &orm{ds: db, lggr: lggr} - - owner := "owner-" + uuid.New().String()[:8] - name := "name-" + uuid.New().String()[:8] - cid := "cid-" + uuid.New().String()[:8] - - t.Run("inserts new spec and new secrets", func(t *testing.T) { //nolint:paralleltest // subtests share database setup - giveURL := "https://example.com/" + uuid.New().String()[:8] - giveBytes, err := crypto.Keccak256([]byte(giveURL)) - require.NoError(t, err) - giveHash := hex.EncodeToString(giveBytes) - giveContent := "some contents" - - spec := &job.WorkflowSpec{ - Workflow: "test_workflow", - Config: "test_config", - WorkflowID: cid, - WorkflowOwner: owner, - WorkflowName: name, - Status: job.WorkflowSpecStatusActive, - BinaryURL: "http://example.com/binary/" + cid, - ConfigURL: "http://example.com/config/" + cid, - CreatedAt: time.Now(), - SpecType: job.WASMFile, - } - - _, err = orm.UpsertWorkflowSpecWithSecrets(ctx, spec, giveURL, giveHash, giveContent) - require.NoError(t, err) - - // Verify the record exists in the database - var dbSpec job.WorkflowSpec - err = db.Get(&dbSpec, `SELECT * FROM workflow_specs WHERE workflow_owner = $1 AND workflow_name = $2`, spec.WorkflowOwner, spec.WorkflowName) - require.NoError(t, err) - require.Equal(t, spec.Workflow, dbSpec.Workflow) - - // Verify the secrets exists in the database - contents, err := orm.GetContents(ctx, giveURL) - require.NoError(t, err) - require.Equal(t, giveContent, contents) - }) - - t.Run("updates existing spec and secrets", func(t *testing.T) { //nolint:paralleltest // subtests share database setup - giveURL := "https://example.com/" + uuid.New().String()[:8] - giveBytes, err := crypto.Keccak256([]byte(giveURL)) - require.NoError(t, err) - giveHash := hex.EncodeToString(giveBytes) - giveContent := "some contents" - - spec := &job.WorkflowSpec{ - Workflow: "test_workflow", - Config: "test_config", - WorkflowID: cid, - WorkflowOwner: owner, - WorkflowName: name, - Status: job.WorkflowSpecStatusActive, - BinaryURL: "http://example.com/binary/" + cid, - ConfigURL: "http://example.com/config/" + cid, - CreatedAt: time.Now(), - SpecType: job.WASMFile, - } - - _, err = orm.UpsertWorkflowSpecWithSecrets(ctx, spec, giveURL, giveHash, giveContent) - require.NoError(t, err) - - // Update the status - spec.Status = job.WorkflowSpecStatusPaused - - _, err = orm.UpsertWorkflowSpecWithSecrets(ctx, spec, giveURL, giveHash, "new contents") - require.NoError(t, err) - - // Verify the record is updated in the database - var dbSpec job.WorkflowSpec - err = db.Get(&dbSpec, `SELECT * FROM workflow_specs WHERE workflow_owner = $1 AND workflow_name = $2`, spec.WorkflowOwner, spec.WorkflowName) - require.NoError(t, err) - require.Equal(t, spec.Config, dbSpec.Config) - - // Verify the secrets is updated in the database - contents, err := orm.GetContents(ctx, giveURL) - require.NoError(t, err) - require.Equal(t, "new contents", contents) - }) - - t.Run("updates existing spec and secrets if spec has executions", func(t *testing.T) { //nolint:paralleltest // subtests share database setup - giveURL := "https://example.com/" + uuid.New().String()[:8] - giveBytes, err := crypto.Keccak256([]byte(giveURL)) - require.NoError(t, err) - giveHash := hex.EncodeToString(giveBytes) - giveContent := "some contents" - - spec := &job.WorkflowSpec{ - Workflow: "test_workflow", - Config: "test_config", - WorkflowID: cid, - WorkflowOwner: owner, - WorkflowName: name, - Status: job.WorkflowSpecStatusActive, - BinaryURL: "http://example.com/binary/" + cid, - ConfigURL: "http://example.com/config/" + cid, - CreatedAt: time.Now(), - SpecType: job.WASMFile, - } - - _, err = orm.UpsertWorkflowSpecWithSecrets(ctx, spec, giveURL, giveHash, giveContent) - require.NoError(t, err) - - _, err = db.ExecContext( - ctx, - `INSERT INTO workflow_executions (id, workflow_id, status, created_at) VALUES ($1, $2, $3, $4)`, - uuid.New().String(), - cid, - "started", - time.Now(), - ) - require.NoError(t, err) - - // Update the status - spec.WorkflowID = "cid-456-" + uuid.New().String()[:8] - - _, err = orm.UpsertWorkflowSpecWithSecrets(ctx, spec, giveURL, giveHash, "new contents") - require.NoError(t, err) - }) -} diff --git a/core/services/workflows/artifacts/store.go b/core/services/workflows/artifacts/store.go deleted file mode 100644 index cfe9d805255..00000000000 --- a/core/services/workflows/artifacts/store.go +++ /dev/null @@ -1,423 +0,0 @@ -package artifacts - -import ( - "context" - "crypto/sha256" - "database/sql" - "encoding/base64" - "encoding/hex" - "encoding/json" - "errors" - "fmt" - "math" - "net/http" - "strings" - "sync" - "time" - - "github.com/jonboulle/clockwork" - - "github.com/smartcontractkit/chainlink-common/keystore/corekeys/workflowkey" - "github.com/smartcontractkit/chainlink-common/pkg/custmsg" - "github.com/smartcontractkit/chainlink-common/pkg/logger" - "github.com/smartcontractkit/chainlink-common/pkg/workflows/secrets" - "github.com/smartcontractkit/chainlink/v2/core/platform" - ghcapabilities "github.com/smartcontractkit/chainlink/v2/core/services/gateway/handlers/capabilities" - "github.com/smartcontractkit/chainlink/v2/core/services/job" - "github.com/smartcontractkit/chainlink/v2/core/services/workflows/types" - "github.com/smartcontractkit/chainlink/v2/core/utils" -) - -type lastFetchedAtMap struct { - m map[string]time.Time - sync.RWMutex -} - -func (l *lastFetchedAtMap) Set(url string, at time.Time) { - l.Lock() - defer l.Unlock() - l.m[url] = at -} - -func (l *lastFetchedAtMap) Get(url string) (time.Time, bool) { - l.RLock() - defer l.RUnlock() - got, ok := l.m[url] - return got, ok -} - -func newLastFetchedAtMap() *lastFetchedAtMap { - return &lastFetchedAtMap{ - m: map[string]time.Time{}, - } -} - -func safeUint32(n uint64) uint32 { - if n > math.MaxUint32 { - return math.MaxUint32 - } - return uint32(n) -} - -type ArtifactConfig struct { - MaxConfigSize uint64 - MaxSecretsSize uint64 - MaxBinarySize uint64 -} - -// By default, if type is unknown, the largest artifact size is 26.4KB. Configure the artifact size -// via the ArtifactConfig to override this default. -const defaultMaxArtifactSizeBytes = 26.4 * utils.KB - -func (cfg *ArtifactConfig) ApplyDefaults() { - if cfg.MaxConfigSize == 0 { - cfg.MaxConfigSize = defaultMaxArtifactSizeBytes - } - if cfg.MaxSecretsSize == 0 { - cfg.MaxSecretsSize = defaultMaxArtifactSizeBytes - } - if cfg.MaxBinarySize == 0 { - cfg.MaxBinarySize = defaultMaxArtifactSizeBytes - } -} - -// logCustMsg emits a custom message to the external sink and logs an error if that fails. -func logCustMsg(ctx context.Context, cma custmsg.MessageEmitter, msg string, log logger.Logger) { - err := cma.Emit(ctx, msg) - if err != nil { - logger.Helper(log, 1).Errorf("failed to send custom message with msg: %s, err: %v", msg, err) - } -} - -var defaultSecretsFreshnessDuration = 24 * time.Hour - -func WithMaxArtifactSize(cfg ArtifactConfig) func(*Store) { - return func(a *Store) { - a.limits = &cfg - } -} - -type SerialisedModuleStore interface { - StoreModule(workflowID string, module []byte) error - GetModulePath(workflowID string) (string, bool, error) - DeleteModule(workflowID string) error -} - -type decryptSecretsFn func(data []byte, owner string) (map[string]string, error) - -type Store struct { - lggr logger.Logger - - // limits sets max artifact sizes to fetch when handling events - limits *ArtifactConfig - - orm WorkflowRegistryDS - - // fetchFn is a function that fetches the contents of a URL with a limit on the size of the response. - fetchFn types.FetcherFunc - - lastFetchedAtMap *lastFetchedAtMap - clock clockwork.Clock - secretsFreshnessDuration time.Duration - - encryptionKey workflowkey.Key - - decryptSecrets decryptSecretsFn - - emitter custmsg.MessageEmitter -} - -func NewStore(lggr logger.Logger, orm WorkflowRegistryDS, fetchFn types.FetcherFunc, clock clockwork.Clock, encryptionKey workflowkey.Key, - emitter custmsg.MessageEmitter, opts ...func(*Store)) *Store { - return NewStoreWithDecryptSecretsFn(lggr, orm, fetchFn, clock, encryptionKey, emitter, - func(data []byte, owner string) (map[string]string, error) { - secretsPayload := secrets.EncryptedSecretsResult{} - err := json.Unmarshal(data, &secretsPayload) - if err != nil { - return nil, fmt.Errorf("could not unmarshal secrets: %w", err) - } - - return secrets.DecryptSecretsForNode(secretsPayload, encryptionKey, owner) - }, - opts...) -} - -func NewStoreWithDecryptSecretsFn(lggr logger.Logger, orm WorkflowRegistryDS, fetchFn types.FetcherFunc, clock clockwork.Clock, encryptionKey workflowkey.Key, - emitter custmsg.MessageEmitter, decryptSecrets decryptSecretsFn, opts ...func(*Store)) *Store { - limits := &ArtifactConfig{} - limits.ApplyDefaults() - - artifactsStore := &Store{ - lggr: lggr, - orm: orm, - fetchFn: fetchFn, - lastFetchedAtMap: newLastFetchedAtMap(), - clock: clock, - limits: limits, - secretsFreshnessDuration: defaultSecretsFreshnessDuration, - encryptionKey: encryptionKey, - emitter: emitter, - decryptSecrets: decryptSecrets, - } - - for _, o := range opts { - o(artifactsStore) - } - - return artifactsStore -} - -// FetchWorkflowArtifacts fetches the workflow spec and config from a cache or the specified URLs if the artifacts have not -// been cached already. Before a workflow can be started this method must be called to ensure all artifacts used by the -// workflow are available from the store. -func (h *Store) FetchWorkflowArtifacts(ctx context.Context, workflowID, binaryURL, configURL string) ([]byte, []byte, error) { - // Check if the workflow spec is already stored in the database - if spec, err := h.orm.GetWorkflowSpecByID(ctx, workflowID); err == nil { - // there is no update in the BinaryURL or ConfigURL, lets decode the stored artifacts - decodedBinary, err := hex.DecodeString(spec.Workflow) - if err != nil { - return nil, nil, fmt.Errorf("failed to decode stored workflow spec: %w", err) - } - return decodedBinary, []byte(spec.Config), nil - } - - // Fetch the binary and config files from the specified URLs. - var ( - binary, decodedBinary, config []byte - err error - ) - - req := ghcapabilities.Request{ - URL: binaryURL, - Method: http.MethodGet, - MaxResponseBytes: safeUint32(h.limits.MaxBinarySize), - WorkflowID: workflowID, - } - binary, err = h.fetchFn(ctx, messageID(binaryURL, workflowID), req) - if err != nil { - return nil, nil, fmt.Errorf("failed to fetch binary from %s : %w", binaryURL, err) - } - - if decodedBinary, err = base64.StdEncoding.DecodeString(string(binary)); err != nil { - return nil, nil, fmt.Errorf("failed to decode binary: %w binary length %d", err, len(binary)) - } - - if configURL != "" { - req := ghcapabilities.Request{ - URL: configURL, - Method: http.MethodGet, - MaxResponseBytes: safeUint32(h.limits.MaxConfigSize), - WorkflowID: workflowID, - } - config, err = h.fetchFn(ctx, messageID(configURL, workflowID), req) - if err != nil { - return nil, nil, fmt.Errorf("failed to fetch config from %s : %w", configURL, err) - } - } - return decodedBinary, config, nil -} - -func (h *Store) GetSecrets(ctx context.Context, secretsURL string, workflowID [32]byte, workflowOwner []byte) ([]byte, error) { - wid := hex.EncodeToString(workflowID[:]) - req := ghcapabilities.Request{ - URL: secretsURL, - Method: http.MethodGet, - MaxResponseBytes: safeUint32(h.limits.MaxSecretsSize), - WorkflowID: wid, - } - fetchedSecrets, fetchErr := h.fetchFn(ctx, messageID(secretsURL, wid), req) - if fetchErr != nil { - return nil, fmt.Errorf("failed to fetch secrets from %s : %w", secretsURL, fetchErr) - } - - return fetchedSecrets, nil -} - -func (h *Store) ValidateSecrets(ctx context.Context, workflowID, workflowOwner string) error { - _, secretsPayload, err := h.orm.GetContentsByWorkflowID(ctx, workflowID) - if err != nil { - // The workflow record was found, but secrets_id was empty. - if errors.Is(err, ErrEmptySecrets) { - return nil - } - - return fmt.Errorf("failed to retrieve secrets by workflow ID: %w", err) - } - - _, decryptErr := h.decryptSecrets([]byte(secretsPayload), workflowOwner) - if decryptErr != nil { - return fmt.Errorf("failed to decrypt secrets: %w", decryptErr) - } - - return nil -} - -func (h *Store) ForceUpdateSecrets( - ctx context.Context, - secretsURLHash []byte, - owner []byte, -) (string, error) { - // Get the URL of the secrets file from the event data - hash := hex.EncodeToString(secretsURLHash) - - url, err := h.orm.GetSecretsURLByHash(ctx, hash) - if err != nil { - return "", fmt.Errorf("failed to get URL by hash %s : %w", hash, err) - } - - ownerHex := hex.EncodeToString(owner) - req := ghcapabilities.Request{ - URL: url, - Method: http.MethodGet, - MaxResponseBytes: safeUint32(h.limits.MaxSecretsSize), - // TODO -- fix, but this is used for rate limiting purposes - WorkflowID: hex.EncodeToString(owner), - } - // Fetch the contents of the secrets file from the url via the fetcher - secrets, err := h.fetchFn(ctx, messageID(url, ownerHex), req) - if err != nil { - return "", err - } - - // Sanity check the payload and ensure we can decrypt it. - // If we can't, let's return an error and we won't store the result in the DB. - _, err = h.decryptSecrets(secrets, hex.EncodeToString(owner)) - if err != nil { - return "", fmt.Errorf("failed to validate secrets: could not decrypt: %w", err) - } - - h.lastFetchedAtMap.Set(hash, h.clock.Now()) - - // Update the secrets in the ORM - if _, err := h.orm.Update(ctx, hash, string(secrets)); err != nil { - return "", fmt.Errorf("failed to update secrets: %w", err) - } - - return string(secrets), nil -} - -func (h *Store) GetSecretsURLByID(ctx context.Context, id int64) (string, error) { - secretsURL, err := h.orm.GetSecretsURLByID(ctx, id) - if err != nil { - return "", fmt.Errorf("failed to get secrets URL by ID: %w", err) - } - return secretsURL, nil -} - -func (h *Store) GetWorkflowSpec(ctx context.Context, workflowOwner string, workflowName string) (*job.WorkflowSpec, error) { - spec, err := h.orm.GetWorkflowSpec(ctx, workflowOwner, workflowName) - return spec, err -} - -func (h *Store) SecretsFor(ctx context.Context, workflowOwner, hexWorkflowName, decodedWorkflowName, workflowID string) (map[string]string, error) { - secretsURLHash, secretsPayload, err := h.orm.GetContentsByWorkflowID(ctx, workflowID) - if err != nil { - // The workflow record was found, but secrets_id was empty. - // Let's just stub out the response. - if errors.Is(err, ErrEmptySecrets) { - return map[string]string{}, nil - } - - return nil, fmt.Errorf("failed to fetch secrets by workflow ID: %w", err) - } - - lastFetchedAt, ok := h.lastFetchedAtMap.Get(secretsURLHash) - if !ok || h.clock.Now().Sub(lastFetchedAt) > h.secretsFreshnessDuration { - updatedSecrets, innerErr := h.refreshSecrets(ctx, workflowOwner, hexWorkflowName, workflowID, secretsURLHash) - if innerErr != nil { - msg := fmt.Sprintf("could not refresh secrets: proceeding with stale secrets for workflowID %s: %s", workflowID, innerErr) - h.lggr.Error(msg) - - logCustMsg( - ctx, - h.emitter.With( - platform.KeyWorkflowID, workflowID, - platform.KeyWorkflowName, decodedWorkflowName, - platform.KeyWorkflowOwner, workflowOwner, - ), - msg, - h.lggr, - ) - } else { - secretsPayload = updatedSecrets - } - } - - return h.decryptSecrets([]byte(secretsPayload), workflowOwner) -} - -func (h *Store) UpsertWorkflowSpec(ctx context.Context, spec *job.WorkflowSpec) (int64, error) { - return h.orm.UpsertWorkflowSpec(ctx, spec) -} - -func (h *Store) UpsertWorkflowSpecWithSecrets(ctx context.Context, entry *job.WorkflowSpec, secretsURL, urlHash, secrets string) (int64, error) { - return h.orm.UpsertWorkflowSpecWithSecrets(ctx, entry, secretsURL, urlHash, secrets) -} - -func (h *Store) GetSecretsURLHash(workflowOwner []byte, secretsURL []byte) ([]byte, error) { - urlHash, err := h.orm.GetSecretsURLHash(workflowOwner, secretsURL) - if err != nil { - return nil, fmt.Errorf("failed to get secrets URL hash: %w", err) - } - return urlHash, nil -} - -// DeleteWorkflowArtifacts removes the workflow spec from the database. If not found, returns nil. -func (h *Store) DeleteWorkflowArtifacts(ctx context.Context, workflowOwner string, workflowName string, workflowID string) error { - err := h.orm.DeleteWorkflowSpec(ctx, workflowOwner, workflowName) - if err != nil { - if errors.Is(err, sql.ErrNoRows) { - h.lggr.Warnw("failed to delete workflow spec: not found", "workflowID", workflowID) - return nil - } - return fmt.Errorf("failed to delete workflow spec: %w", err) - } - - return nil -} - -func (h *Store) GetWasmBinary(ctx context.Context, workflowID string) ([]byte, error) { - spec, err := h.orm.GetWorkflowSpecByID(ctx, workflowID) - if err != nil { - return nil, fmt.Errorf("failed to get workflow spec by workflow ID: %w", err) - } - - // there is no update in the BinaryURL or ConfigURL, lets decode the stored artifacts - decodedBinary, err := hex.DecodeString(spec.Workflow) - if err != nil { - return nil, fmt.Errorf("failed to decode stored workflow string: %w", err) - } - - return decodedBinary, nil -} - -func (h *Store) refreshSecrets(ctx context.Context, workflowOwner, workflowName, workflowID, secretsURLHash string) (string, error) { - owner, err := hex.DecodeString(workflowOwner) - if err != nil { - return "", err - } - - decodedHash, err := hex.DecodeString(secretsURLHash) - if err != nil { - return "", err - } - - updatedSecrets, err := h.ForceUpdateSecrets( - ctx, decodedHash, owner) - if err != nil { - return "", err - } - - return updatedSecrets, nil -} - -func messageID(url string, parts ...string) string { - h := sha256.New() - h.Write([]byte(url)) - for _, p := range parts { - h.Write([]byte(p)) - } - hash := hex.EncodeToString(h.Sum(nil)) - p := []string{ghcapabilities.MethodWorkflowSyncer, hash} - return strings.Join(p, "/") -} diff --git a/core/services/workflows/artifacts/store_test.go b/core/services/workflows/artifacts/store_test.go deleted file mode 100644 index bc2fe7faa78..00000000000 --- a/core/services/workflows/artifacts/store_test.go +++ /dev/null @@ -1,310 +0,0 @@ -package artifacts - -import ( - "context" - "database/sql" - "encoding/hex" - "encoding/json" - "errors" - "testing" - "time" - - "github.com/jonboulle/clockwork" - "github.com/stretchr/testify/assert" - "github.com/stretchr/testify/require" - - "github.com/smartcontractkit/chainlink-common/keystore/corekeys/workflowkey" - "github.com/smartcontractkit/chainlink-common/pkg/custmsg" - "github.com/smartcontractkit/chainlink-common/pkg/workflows/secrets" - "github.com/smartcontractkit/chainlink/v2/core/internal/testutils/pgtest" - "github.com/smartcontractkit/chainlink/v2/core/logger" - ghcapabilities "github.com/smartcontractkit/chainlink/v2/core/services/gateway/handlers/capabilities" - "github.com/smartcontractkit/chainlink/v2/core/services/job" -) - -func Test_Handler_SecretsFor(t *testing.T) { - t.Parallel() - lggr := logger.TestLogger(t) - db := pgtest.NewSqlxDB(t) - orm := &orm{ds: db, lggr: lggr} - - workflowOwner := hex.EncodeToString([]byte("anOwner")) - workflowName := "aName" - workflowID := "anID" - decodedWorkflowName := "decodedName" - encryptionKey, err := workflowkey.New() - require.NoError(t, err) - - url := "http://example.com" - hash := hex.EncodeToString([]byte(url)) - secretsPayload, err := generateSecrets(workflowOwner, map[string][]string{"Foo": {"Bar"}}, encryptionKey) - require.NoError(t, err) - secretsID, err := orm.Create(t.Context(), url, hash, string(secretsPayload)) - require.NoError(t, err) - - _, err = orm.UpsertWorkflowSpec(t.Context(), &job.WorkflowSpec{ - Workflow: "", - Config: "", - SecretsID: sql.NullInt64{Int64: secretsID, Valid: true}, - WorkflowID: workflowID, - WorkflowOwner: workflowOwner, - WorkflowName: workflowName, - BinaryURL: "", - ConfigURL: "", - CreatedAt: time.Now(), - SpecType: job.DefaultSpecType, - }) - require.NoError(t, err) - - fetcher := &mockFetcher{ - responseMap: map[string]mockFetchResp{ - url: {Err: errors.New("could not fetch")}, - }, - } - - require.NoError(t, err) - h := NewStore( - lggr, - orm, - fetcher.Fetch, - clockwork.NewFakeClock(), - encryptionKey, - custmsg.NewLabeler(), - ) - expectedSecrets := map[string]string{ - "Foo": "Bar", - } - decrypter := newMockDecrypter() - decrypter.registerMock(secretsPayload, workflowOwner, expectedSecrets, nil) - h.decryptSecrets = decrypter.decryptSecrets - gotSecrets, err := h.SecretsFor(t.Context(), workflowOwner, workflowName, decodedWorkflowName, workflowID) - require.NoError(t, err) - - assert.Equal(t, expectedSecrets, gotSecrets) -} - -func Test_Handler_SecretsFor_RefreshesSecrets(t *testing.T) { - t.Parallel() - lggr := logger.TestLogger(t) - db := pgtest.NewSqlxDB(t) - orm := &orm{ds: db, lggr: lggr} - - workflowOwner := hex.EncodeToString([]byte("anOwner")) - workflowName := "aName" - workflowID := "anID" - decodedWorkflowName := "decodedName" - encryptionKey, err := workflowkey.New() - require.NoError(t, err) - - secretsPayload, err := generateSecrets(workflowOwner, map[string][]string{"Foo": {"Bar"}}, encryptionKey) - require.NoError(t, err) - - url := "http://example.com" - hash := hex.EncodeToString([]byte(url)) - - secretsID, err := orm.Create(t.Context(), url, hash, string(secretsPayload)) - require.NoError(t, err) - - _, err = orm.UpsertWorkflowSpec(t.Context(), &job.WorkflowSpec{ - Workflow: "", - Config: "", - SecretsID: sql.NullInt64{Int64: secretsID, Valid: true}, - WorkflowID: workflowID, - WorkflowOwner: workflowOwner, - WorkflowName: workflowName, - BinaryURL: "", - ConfigURL: "", - CreatedAt: time.Now(), - SpecType: job.DefaultSpecType, - }) - require.NoError(t, err) - - secretsPayload, err = generateSecrets(workflowOwner, map[string][]string{"Baz": {"Bar"}}, encryptionKey) - require.NoError(t, err) - fetcher := &mockFetcher{ - responseMap: map[string]mockFetchResp{ - url: {Body: secretsPayload}, - }, - } - - h := NewStore( - lggr, - orm, - fetcher.Fetch, - clockwork.NewFakeClock(), - encryptionKey, - custmsg.NewLabeler(), - ) - - expectedSecrets := map[string]string{ - "Baz": "Bar", - } - decrypter := newMockDecrypter() - decrypter.registerMock(secretsPayload, workflowOwner, expectedSecrets, nil) - h.decryptSecrets = decrypter.decryptSecrets - - gotSecrets, err := h.SecretsFor(t.Context(), workflowOwner, workflowName, decodedWorkflowName, workflowID) - require.NoError(t, err) - - assert.Equal(t, expectedSecrets, gotSecrets) -} - -func Test_Handler_SecretsFor_RefreshLogic(t *testing.T) { - t.Parallel() - lggr := logger.TestLogger(t) - db := pgtest.NewSqlxDB(t) - orm := &orm{ds: db, lggr: lggr} - - workflowOwner := hex.EncodeToString([]byte("anOwner")) - workflowName := "aName" - workflowID := "anID" - decodedWorkflowName := "decodedName" - encryptionKey, err := workflowkey.New() - require.NoError(t, err) - - secretsPayload, err := generateSecrets(workflowOwner, map[string][]string{"Foo": {"Bar"}}, encryptionKey) - require.NoError(t, err) - - url := "http://example.com" - hash := hex.EncodeToString([]byte(url)) - - secretsID, err := orm.Create(t.Context(), url, hash, string(secretsPayload)) - require.NoError(t, err) - - _, err = orm.UpsertWorkflowSpec(t.Context(), &job.WorkflowSpec{ - Workflow: "", - Config: "", - SecretsID: sql.NullInt64{Int64: secretsID, Valid: true}, - WorkflowID: workflowID, - WorkflowOwner: workflowOwner, - WorkflowName: workflowName, - BinaryURL: "", - ConfigURL: "", - CreatedAt: time.Now(), - SpecType: job.DefaultSpecType, - }) - require.NoError(t, err) - - fetcher := &mockFetcher{ - responseMap: map[string]mockFetchResp{ - url: { - Body: secretsPayload, - }, - }, - } - clock := clockwork.NewFakeClock() - - h := NewStore( - lggr, - orm, - fetcher.Fetch, - clock, - encryptionKey, - custmsg.NewLabeler(), - ) - - expectedSecrets := map[string]string{ - "Foo": "Bar", - } - decrypter := newMockDecrypter() - decrypter.registerMock(secretsPayload, workflowOwner, expectedSecrets, nil) - h.decryptSecrets = decrypter.decryptSecrets - - gotSecrets, err := h.SecretsFor(t.Context(), workflowOwner, workflowName, decodedWorkflowName, workflowID) - require.NoError(t, err) - - assert.Equal(t, expectedSecrets, gotSecrets) - - // Now stub out an unparseable response, since we already fetched it recently above, we shouldn't need to refetch - // SecretsFor should still succeed. - fetcher.responseMap[url] = mockFetchResp{} - - gotSecrets, err = h.SecretsFor(t.Context(), workflowOwner, workflowName, decodedWorkflowName, workflowID) - require.NoError(t, err) - - assert.Equal(t, expectedSecrets, gotSecrets) - - secretsPayload, err = generateSecrets(workflowOwner, map[string][]string{"Baz": {"Bar"}}, encryptionKey) - require.NoError(t, err) - fetcher.responseMap[url] = mockFetchResp{ - Body: secretsPayload, - } - - expectedSecrets = map[string]string{ - "Baz": "Bar", - } - decrypter.registerMock(secretsPayload, workflowOwner, expectedSecrets, nil) - - // Now advance so that we hit the freshness limit - clock.Advance(48 * time.Hour) - - gotSecrets, err = h.SecretsFor(t.Context(), workflowOwner, workflowName, decodedWorkflowName, workflowID) - require.NoError(t, err) - assert.Equal(t, expectedSecrets, gotSecrets) -} - -func generateSecrets(workflowOwner string, secretsMap map[string][]string, encryptionKey workflowkey.Key) ([]byte, error) { - sm, secretsEnvVars, err := secrets.EncryptSecretsForNodes( - workflowOwner, - secretsMap, - map[string][32]byte{ - "p2pId": encryptionKey.PublicKey(), - }, - secrets.SecretsConfig{}, - ) - if err != nil { - return nil, err - } - return json.Marshal(secrets.EncryptedSecretsResult{ - EncryptedSecrets: sm, - Metadata: secrets.Metadata{ - WorkflowOwner: workflowOwner, - EnvVarsAssignedToNodes: secretsEnvVars, - NodePublicEncryptionKeys: map[string]string{ - "p2pId": encryptionKey.PublicKeyString(), - }, - }, - }) -} - -type mockFetchResp struct { - Body []byte - Err error -} - -type mockFetcher struct { - responseMap map[string]mockFetchResp -} - -func (m *mockFetcher) Fetch(_ context.Context, mid string, req ghcapabilities.Request) ([]byte, error) { - return m.responseMap[req.URL].Body, m.responseMap[req.URL].Err -} - -type decryptSecretsOutput struct { - output map[string]string - err error -} - -type mockDecrypter struct { - mocks map[string]decryptSecretsOutput -} - -func (m *mockDecrypter) decryptSecrets(data []byte, owner string) (map[string]string, error) { - input := string(data) + owner - mock, exists := m.mocks[input] - if exists { - return mock.output, mock.err - } - return map[string]string{}, nil -} - -func (m *mockDecrypter) registerMock(data []byte, owner string, output map[string]string, err error) { - input := string(data) + owner - m.mocks[input] = decryptSecretsOutput{output: output, err: err} -} - -func newMockDecrypter() *mockDecrypter { - return &mockDecrypter{ - mocks: map[string]decryptSecretsOutput{}, - } -} diff --git a/core/services/workflows/artifacts/v2/store_test.go b/core/services/workflows/artifacts/v2/store_test.go index ebd3aec3660..0f1916ef6c0 100644 --- a/core/services/workflows/artifacts/v2/store_test.go +++ b/core/services/workflows/artifacts/v2/store_test.go @@ -63,7 +63,7 @@ func Test_Store_DeleteWorkflowArtifacts(t *testing.T) { BinaryURL: "", ConfigURL: "", CreatedAt: time.Now(), - SpecType: job.DefaultSpecType, + SpecType: job.WASMFile, }) require.NoError(t, err) @@ -117,7 +117,7 @@ func Test_Store_DeleteWorkflowArtifactsBatch(t *testing.T) { WorkflowOwner: owner, WorkflowName: "name-" + id, CreatedAt: time.Now(), - SpecType: job.DefaultSpecType, + SpecType: job.WASMFile, }) require.NoError(t, err) } diff --git a/core/services/workflows/cmd/cre/README.md b/core/services/workflows/cmd/cre/README.md index 1a8236cbeeb..e21052f8855 100644 --- a/core/services/workflows/cmd/cre/README.md +++ b/core/services/workflows/cmd/cre/README.md @@ -16,21 +16,6 @@ make install-plugins-private ``` Run `$GOBIN/cron -h` to confirm the installation. -### Legacy `data_feeds` Example - -1. Build the workflow: - -```bash -cd core/services/workflows/cmd/cre -GOOS=wasip1 GOARCH=wasm go build -o data_feeds.wasm ./examples/legacy/data_feeds/data_feeds_workflow.go -``` - -2. Run the engine with the workflow: - -```bash -go run . --wasm data_feeds.wasm --config ./examples/legacy/data_feeds/config_10_feeds.json 2> stderr.log -``` - ### V2 `cron` Example ("No DAG") Requires that the `cron` capability be installed on the `$GOBIN` path. See [here](#installing-capability-binaries). diff --git a/system-tests/lib/cre/don/jobs/jobs.go b/system-tests/lib/cre/don/jobs/jobs.go index d90f30d0fa0..98ec7db1784 100644 --- a/system-tests/lib/cre/don/jobs/jobs.go +++ b/system-tests/lib/cre/don/jobs/jobs.go @@ -142,10 +142,10 @@ func accept(ctx context.Context, node *cre.Node, proposalID, jobSpec string) err err = approveJobProposalSpec(ctx, node, proposalID) } if err != nil { - // Workflow and CRE settings specs get auto-approved by the node on proposal, so a + // CRE settings specs get auto-approved by the node on proposal, so a // subsequent explicit approve races into an already-approved spec — tolerate that. if strings.Contains(err.Error(), "cannot approve an approved spec") && - (strings.Contains(jobSpec, `type = "workflow"`) || strings.Contains(jobSpec, `type = "cresettings"`)) { + strings.Contains(jobSpec, `type = "cresettings"`) { return nil } fmt.Println("Failed jobspec proposal for node ", node.Name) diff --git a/system-tests/lib/cre/don/jobs/jobs_test.go b/system-tests/lib/cre/don/jobs/jobs_test.go index 00c17e880d2..9a4d0aa314b 100644 --- a/system-tests/lib/cre/don/jobs/jobs_test.go +++ b/system-tests/lib/cre/don/jobs/jobs_test.go @@ -136,7 +136,9 @@ func TestApproveMissingProposalMatch(t *testing.T) { require.ErrorContains(t, err, "no job proposal found for job spec missing-spec") } -func TestAcceptTreatsApprovedWorkflowSpecAsSuccess(t *testing.T) { +func TestAcceptAutoApprovedSpecTolerance(t *testing.T) { + t.Parallel() + restoreApprove := approveJobProposalSpec t.Cleanup(func() { approveJobProposalSpec = restoreApprove @@ -146,6 +148,13 @@ func TestAcceptTreatsApprovedWorkflowSpecAsSuccess(t *testing.T) { return errors.New("cannot approve an approved spec") } - err := accept(context.Background(), &cre.Node{Name: "node-a"}, "proposal-id", `type = "workflow"`) + // CRE settings specs get auto-approved by the node on proposal, so a + // subsequent explicit approve that races into an already-approved spec is tolerated. + err := accept(context.Background(), &cre.Node{Name: "node-a"}, "proposal-id", `type = "cresettings"`) require.NoError(t, err) + + // Workflow specs are no longer auto-approved, so the same race is not tolerated. + err = accept(context.Background(), &cre.Node{Name: "node-a"}, "proposal-id", `type = "workflow"`) + require.Error(t, err) + require.ErrorContains(t, err, "failed to accept job for node node-a") }