-
Notifications
You must be signed in to change notification settings - Fork 2k
Common Job spec reporting via Beholder #23227
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Draft
skudasov
wants to merge
3
commits into
develop
Choose a base branch
from
RANE-4655-common-job-spec-metrics
base: develop
Could not load branches
Branch not found: {{ refName }}
Loading
Could not load tags
Nothing to show
Loading
Are you sure you want to change the base?
Some commits from the old base branch may be removed from the timeline,
and old review comments may become outdated.
+422
−15
Draft
Changes from all commits
Commits
Show all changes
3 commits
Select commit
Hold shift + click to select a range
File filter
Filter by extension
Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
There are no files selected for viewing
141 changes: 141 additions & 0 deletions
141
core/services/nodestatusreporter/jobspec/cl_job_info.go
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,141 @@ | ||
| package jobspec | ||
|
|
||
| import ( | ||
| "context" | ||
| "fmt" | ||
| "time" | ||
|
|
||
| "github.com/pelletier/go-toml" | ||
| "google.golang.org/protobuf/proto" | ||
| "google.golang.org/protobuf/types/known/timestamppb" | ||
|
|
||
| "github.com/smartcontractkit/chainlink-common/pkg/beholder" | ||
| commonv1 "github.com/smartcontractkit/chainlink-protos/node-platform/common/v1" | ||
|
|
||
| "github.com/smartcontractkit/chainlink/v2/core/services/job" | ||
| ) | ||
|
|
||
| // CLJobInfo is the generic, job-type-agnostic half of this reporter. Where | ||
| // JobSpecEvent models one job type (OCR2) field by field, CLJobInfo carries the | ||
| // job's common identity plus the complete definition as a raw TOML string, so | ||
| // every job the node runs is reported through a single schema with no | ||
| // per-type code here and none needed for future job types. | ||
| // | ||
| // It is emitted on the same triggers and from the same service as JobSpecEvent | ||
| // rather than from a parallel one, so there is exactly one place in the node | ||
| // that reports what jobs it runs. Once consumers have migrated, the OCR2-only | ||
| // half can be deleted from here without touching the wiring. | ||
| const ( | ||
| // Domain, Entity and DataSchema identify CLJobInfo telemetry on Beholder. | ||
| Domain = "node-platform" | ||
| Entity = "common.v1.CLJobInfo" | ||
| DataSchema = "/node-platform/common/v1" | ||
| ) | ||
|
|
||
| // NodeIdentity is the node-level context attached to every emitted CLJobInfo. | ||
| type NodeIdentity struct { | ||
| CSAPublicKey string | ||
| NodeVersion string | ||
| Hostname string | ||
| } | ||
|
|
||
| // JobProposal is the Job Distributor provenance for a job that arrived as an | ||
| // approved job proposal. Jobs created directly (CLI, UI, TOML on disk) have no | ||
| // proposal, and the zero value leaves the corresponding CLJobInfo fields unset | ||
| // — which is how a consumer tells a managed job from an unmanaged one. | ||
| type JobProposal struct { | ||
| FeedsManagerID int64 | ||
| RemoteUUID string | ||
| SpecVersion int32 | ||
| ProposedAt time.Time | ||
| ApprovedAt time.Time | ||
| } | ||
|
|
||
| // BuildCLJobInfo converts any job.Job into its generic CLJobInfo representation. | ||
| // | ||
| // prop is optional: pass nil for a job with no Job Distributor proposal. | ||
| // | ||
| // If the job cannot be TOML-encoded, BuildCLJobInfo still returns a fully | ||
| // populated identity payload (with an empty SpecToml) alongside the encoding | ||
| // error, so callers can choose to emit the envelope and log the failure rather | ||
| // than drop the event. | ||
| func BuildCLJobInfo(jb job.Job, trigger commonv1.CLJobInfoTrigger, id NodeIdentity, prop *JobProposal, now time.Time) (*commonv1.CLJobInfo, error) { | ||
| info := &commonv1.CLJobInfo{ | ||
| CsaPublicKey: id.CSAPublicKey, | ||
| NodeVersion: id.NodeVersion, | ||
| Hostname: id.Hostname, | ||
| ExternalJobId: jb.ExternalJobID.String(), | ||
| JobId: jb.ID, | ||
| Name: jb.Name.ValueOrZero(), | ||
| JobType: string(jb.Type), | ||
| SchemaVersion: jb.SchemaVersion, | ||
| ForwardingAllowed: jb.ForwardingAllowed, | ||
| CreatedAt: timestampOrNil(jb.CreatedAt), | ||
| Trigger: trigger, | ||
| Timestamp: timestamppb.New(now), | ||
| } | ||
| if jb.GasLimit.Valid { | ||
| info.GasLimit = new(jb.GasLimit.Uint32) | ||
| } | ||
| if jb.StreamID != nil { | ||
| info.StreamId = new(*jb.StreamID) | ||
| } | ||
| if prop != nil { | ||
| info.FeedsManagerId = &prop.FeedsManagerID | ||
| info.RemoteUuid = &prop.RemoteUUID | ||
| info.SpecVersion = &prop.SpecVersion | ||
| info.ProposedAt = timestampOrNil(prop.ProposedAt) | ||
| info.ApprovedAt = timestampOrNil(prop.ApprovedAt) | ||
| } | ||
|
|
||
| specTOML, err := jobTOML(jb) | ||
| if err != nil { | ||
| return info, fmt.Errorf("encoding job %s (%d) spec to TOML: %w", jb.ExternalJobID, jb.ID, err) | ||
| } | ||
| info.SpecToml = specTOML | ||
|
|
||
| return info, nil | ||
| } | ||
|
|
||
| // EmitCLJobInfo marshals a CLJobInfo and publishes it to Beholder. | ||
| func EmitCLJobInfo(ctx context.Context, emitter beholder.Emitter, info *commonv1.CLJobInfo) error { | ||
| payload, err := proto.Marshal(info) | ||
| if err != nil { | ||
| return fmt.Errorf("marshaling CLJobInfo: %w", err) | ||
| } | ||
|
|
||
| err = emitter.Emit(ctx, payload, | ||
| beholder.AttrKeyDomain, Domain, | ||
| beholder.AttrKeyEntity, Entity, | ||
| beholder.AttrKeyDataSchema, DataSchema, | ||
| ) | ||
| if err != nil { | ||
| return fmt.Errorf("emitting CLJobInfo: %w", err) | ||
| } | ||
| return nil | ||
| } | ||
|
|
||
| // jobTOML serializes the entire job definition to TOML. Marshaling the whole | ||
| // job.Job captures both the common top-level fields and the single active | ||
| // type-specific spec, so all fields for any job type are included without | ||
| // enumerating them. | ||
| func jobTOML(jb job.Job) (string, error) { | ||
| out, err := toml.Marshal(jb) | ||
| if err != nil { | ||
| return "", err | ||
| } | ||
| return string(out), nil | ||
| } | ||
|
|
||
| // timestampOrNil converts t to a protobuf Timestamp, leaving an unset time as | ||
| // nil rather than mapping it onto the epoch. google.protobuf.Timestamp is used | ||
| // throughout the Job Distributor protos and, unlike an RFC3339Nano string, | ||
| // orders correctly for consumers: Go trims trailing zeros from the fractional | ||
| // seconds, so those strings are variable-width and do not sort lexicographically | ||
| // in chronological order. | ||
| func timestampOrNil(t time.Time) *timestamppb.Timestamp { | ||
| if t.IsZero() { | ||
| return nil | ||
| } | ||
| return timestamppb.New(t) | ||
| } | ||
196 changes: 196 additions & 0 deletions
196
core/services/nodestatusreporter/jobspec/cl_job_info_test.go
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,196 @@ | ||
| package jobspec_test | ||
|
|
||
| import ( | ||
| "testing" | ||
| "time" | ||
|
|
||
| "github.com/pelletier/go-toml" | ||
| "github.com/stretchr/testify/require" | ||
| "google.golang.org/protobuf/proto" | ||
| "gopkg.in/guregu/null.v4" | ||
|
|
||
| "github.com/smartcontractkit/chainlink-common/pkg/beholder" | ||
| "github.com/smartcontractkit/chainlink-common/pkg/beholder/beholdertest" | ||
| "github.com/smartcontractkit/chainlink-common/pkg/sqlutil" | ||
| commontypes "github.com/smartcontractkit/chainlink-common/pkg/types" | ||
| evmtypes "github.com/smartcontractkit/chainlink-evm/pkg/types" | ||
| commonv1 "github.com/smartcontractkit/chainlink-protos/node-platform/common/v1" | ||
|
|
||
| "github.com/smartcontractkit/chainlink/v2/core/services/job" | ||
| "github.com/smartcontractkit/chainlink/v2/core/services/nodestatusreporter/jobspec" | ||
| "github.com/smartcontractkit/chainlink/v2/core/services/pipeline" | ||
| ) | ||
|
|
||
| func clJobInfoSampleJob() job.Job { | ||
| streamID := uint32(42) | ||
| return job.Job{ | ||
| ID: 7, | ||
| Name: null.StringFrom("my-ocr2-job"), | ||
| Type: job.OffchainReporting2, | ||
| SchemaVersion: 1, | ||
| ForwardingAllowed: true, | ||
| StreamID: &streamID, | ||
| CreatedAt: time.Date(2026, 7, 24, 10, 0, 0, 0, time.UTC), | ||
| OCR2OracleSpec: &job.OCR2OracleSpec{ | ||
| Relay: "evm", | ||
| ChainID: "1", | ||
| PluginType: commontypes.Median, | ||
| ContractID: "0xcccccccccccccccccccccccccccccccccccccccc", | ||
| TransmitterID: null.StringFrom("0x1111111111111111111111111111111111111111"), | ||
| RelayConfig: job.JSONConfig{ | ||
| "chainID": "1", | ||
| "sendingKeys": []any{"0x1111111111111111111111111111111111111111"}, | ||
| }, | ||
| }, | ||
| Pipeline: pipeline.Pipeline{Tasks: []pipeline.Task{ | ||
| &pipeline.ETHTxTask{From: "0xaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa"}, | ||
| }}, | ||
| } | ||
| } | ||
|
|
||
| // TestBuildCLJobInfo_EncodesFullSpecAsTOML is the load-bearing check: an | ||
| // arbitrary job must round-trip to TOML with no per-type code. | ||
| func TestBuildCLJobInfo_EncodesFullSpecAsTOML(t *testing.T) { | ||
| jb := clJobInfoSampleJob() | ||
| id := jobspec.NodeIdentity{CSAPublicKey: "csa", NodeVersion: "1.2.3", Hostname: "host-1"} | ||
|
|
||
| info, err := jobspec.BuildCLJobInfo(jb, commonv1.CLJobInfoTrigger_CL_JOB_INFO_TRIGGER_CREATE, id, nil, time.Date(2026, 7, 24, 12, 0, 0, 0, time.UTC)) | ||
| require.NoError(t, err) | ||
|
|
||
| require.Equal(t, "csa", info.CsaPublicKey) | ||
| require.Equal(t, "1.2.3", info.NodeVersion) | ||
| require.Equal(t, "host-1", info.Hostname) | ||
| require.Equal(t, int32(7), info.JobId) | ||
| require.Equal(t, "my-ocr2-job", info.Name) | ||
| require.Equal(t, "offchainreporting2", info.JobType) | ||
| require.Equal(t, uint32(1), info.SchemaVersion) | ||
| require.True(t, info.ForwardingAllowed) | ||
| require.NotNil(t, info.StreamId) | ||
| require.Equal(t, uint32(42), *info.StreamId) | ||
| require.Equal(t, commonv1.CLJobInfoTrigger_CL_JOB_INFO_TRIGGER_CREATE, info.Trigger) | ||
| require.NotNil(t, info.Timestamp) | ||
| require.Equal(t, time.Date(2026, 7, 24, 12, 0, 0, 0, time.UTC), info.Timestamp.AsTime()) | ||
| require.NotNil(t, info.CreatedAt) | ||
| require.Equal(t, time.Date(2026, 7, 24, 10, 0, 0, 0, time.UTC), info.CreatedAt.AsTime()) | ||
|
|
||
| // spec_toml must be valid TOML and contain type-specific spec data. | ||
| require.NotEmpty(t, info.SpecToml) | ||
| var decoded map[string]any | ||
| require.NoError(t, toml.Unmarshal([]byte(info.SpecToml), &decoded)) | ||
| require.Contains(t, info.SpecToml, "median") | ||
| require.Contains(t, info.SpecToml, "0xcccccccccccccccccccccccccccccccccccccccc") | ||
| } | ||
|
|
||
| func TestBuildCLJobInfo_HandlesMultipleJobTypesGenerically(t *testing.T) { | ||
| jobs := []job.Job{ | ||
| {Type: job.VRF, VRFSpec: &job.VRFSpec{ | ||
| EVMChainID: sqlutil.NewI(4), | ||
| FromAddresses: []evmtypes.EIP55Address{evmtypes.MustEIP55Address("0x6666666666666666666666666666666666666666")}, | ||
| }}, | ||
| {Type: job.BlockhashStore, BlockhashStoreSpec: &job.BlockhashStoreSpec{EVMChainID: sqlutil.NewI(5)}}, | ||
| } | ||
| for _, jb := range jobs { | ||
| info, err := jobspec.BuildCLJobInfo(jb, commonv1.CLJobInfoTrigger_CL_JOB_INFO_TRIGGER_HEARTBEAT, jobspec.NodeIdentity{}, nil, time.Now()) | ||
| require.NoErrorf(t, err, "job type %s should encode without per-type code", jb.Type) | ||
| require.NotEmpty(t, info.SpecToml) | ||
| } | ||
| } | ||
|
|
||
| // TestBuildCLJobInfo_CarriesJobDistributorProvenance covers the JD join key: | ||
| // remote_uuid is what links this event back to api.job.v1.Job.uuid. | ||
| func TestBuildCLJobInfo_CarriesJobDistributorProvenance(t *testing.T) { | ||
| proposedAt := time.Date(2026, 7, 20, 9, 0, 0, 0, time.UTC) | ||
| approvedAt := time.Date(2026, 7, 24, 10, 0, 0, 0, time.UTC) | ||
| prop := &jobspec.JobProposal{ | ||
| FeedsManagerID: 3, | ||
| RemoteUUID: "6d7d9d1a-0d0f-4b3f-9a2f-2e4a1c0b8d55", | ||
| SpecVersion: 2, | ||
| ProposedAt: proposedAt, | ||
| ApprovedAt: approvedAt, | ||
| } | ||
|
|
||
| info, err := jobspec.BuildCLJobInfo(clJobInfoSampleJob(), commonv1.CLJobInfoTrigger_CL_JOB_INFO_TRIGGER_CREATE, jobspec.NodeIdentity{}, prop, time.Now()) | ||
| require.NoError(t, err) | ||
|
|
||
| require.NotNil(t, info.FeedsManagerId) | ||
| require.Equal(t, int64(3), *info.FeedsManagerId) | ||
| require.NotNil(t, info.RemoteUuid) | ||
| require.Equal(t, "6d7d9d1a-0d0f-4b3f-9a2f-2e4a1c0b8d55", *info.RemoteUuid) | ||
| require.NotNil(t, info.SpecVersion) | ||
| require.Equal(t, int32(2), *info.SpecVersion) | ||
| require.NotNil(t, info.ProposedAt) | ||
| require.Equal(t, proposedAt, info.ProposedAt.AsTime()) | ||
| require.NotNil(t, info.ApprovedAt) | ||
| require.Equal(t, approvedAt, info.ApprovedAt.AsTime()) | ||
| } | ||
|
|
||
| // TestBuildCLJobInfo_UnmanagedJobHasNoProvenance: an unset feeds_manager_id is | ||
| // how a consumer tells a directly-created job from a JD-managed one. | ||
| func TestBuildCLJobInfo_UnmanagedJobHasNoProvenance(t *testing.T) { | ||
| info, err := jobspec.BuildCLJobInfo(clJobInfoSampleJob(), commonv1.CLJobInfoTrigger_CL_JOB_INFO_TRIGGER_CREATE, jobspec.NodeIdentity{}, nil, time.Now()) | ||
| require.NoError(t, err) | ||
|
|
||
| require.Nil(t, info.FeedsManagerId) | ||
| require.Nil(t, info.RemoteUuid) | ||
| require.Nil(t, info.SpecVersion) | ||
| require.Nil(t, info.ProposedAt) | ||
| require.Nil(t, info.ApprovedAt) | ||
| } | ||
|
|
||
| func TestEmitCLJobInfo_PublishesToBeholder(t *testing.T) { | ||
| obs := beholdertest.NewObserver(t) | ||
|
|
||
| info, err := jobspec.BuildCLJobInfo(clJobInfoSampleJob(), commonv1.CLJobInfoTrigger_CL_JOB_INFO_TRIGGER_CREATE, jobspec.NodeIdentity{CSAPublicKey: "csa"}, nil, time.Now()) | ||
| require.NoError(t, err) | ||
| require.NoError(t, jobspec.EmitCLJobInfo(t.Context(), beholder.GetEmitter(), info)) | ||
|
|
||
| msgs := obs.Messages(t, beholder.AttrKeyEntity, jobspec.Entity) | ||
| require.NotEmpty(t, msgs) | ||
|
|
||
| msg := msgs[0] | ||
| require.Equal(t, jobspec.Domain, msg.Attrs[beholder.AttrKeyDomain]) | ||
| require.Equal(t, jobspec.DataSchema, msg.Attrs[beholder.AttrKeyDataSchema]) | ||
|
|
||
| var payload commonv1.CLJobInfo | ||
| require.NoError(t, proto.Unmarshal(msg.Body, &payload)) | ||
| require.Equal(t, "csa", payload.CsaPublicKey) | ||
| require.Equal(t, "offchainreporting2", payload.JobType) | ||
| require.NotEmpty(t, payload.SpecToml) | ||
| } | ||
|
|
||
| // TestBuildCLJobInfo_TimestampsRoundTripExactly guards the reason these fields | ||
| // are google.protobuf.Timestamp rather than RFC3339Nano strings: Go trims | ||
| // trailing zeros from the fractional seconds, so string-encoded times are | ||
| // variable-width and do not sort lexicographically in chronological order — a | ||
| // whole-second time sorts after every sub-second one in the same second. | ||
| func TestBuildCLJobInfo_TimestampsRoundTripExactly(t *testing.T) { | ||
| for _, tc := range []struct { | ||
| name string | ||
| at time.Time | ||
| }{ | ||
| {"whole second", time.Date(2026, 7, 24, 10, 0, 0, 0, time.UTC)}, | ||
| {"tenth of a second", time.Date(2026, 7, 24, 10, 0, 0, 100000000, time.UTC)}, | ||
| {"sub-millisecond", time.Date(2026, 7, 24, 10, 0, 0, 123400000, time.UTC)}, | ||
| {"nanosecond", time.Date(2026, 7, 24, 10, 0, 0, 123456789, time.UTC)}, | ||
| } { | ||
| t.Run(tc.name, func(t *testing.T) { | ||
| jb := clJobInfoSampleJob() | ||
| jb.CreatedAt = tc.at | ||
|
|
||
| info, err := jobspec.BuildCLJobInfo(jb, commonv1.CLJobInfoTrigger_CL_JOB_INFO_TRIGGER_HEARTBEAT, jobspec.NodeIdentity{}, nil, time.Now()) | ||
| require.NoError(t, err) | ||
| require.NotNil(t, info.CreatedAt) | ||
| require.Equal(t, tc.at, info.CreatedAt.AsTime()) | ||
| }) | ||
| } | ||
| } | ||
|
|
||
| // TestBuildCLJobInfo_ZeroTimeIsUnset: an absent time must be nil, not the epoch. | ||
| func TestBuildCLJobInfo_ZeroTimeIsUnset(t *testing.T) { | ||
| jb := clJobInfoSampleJob() | ||
| jb.CreatedAt = time.Time{} | ||
|
|
||
| info, err := jobspec.BuildCLJobInfo(jb, commonv1.CLJobInfoTrigger_CL_JOB_INFO_TRIGGER_HEARTBEAT, jobspec.NodeIdentity{}, nil, time.Now()) | ||
| require.NoError(t, err) | ||
| require.Nil(t, info.CreatedAt) | ||
| } |
Oops, something went wrong.
Add this suggestion to a batch that can be applied as a single commit.
This suggestion is invalid because no changes were made to the code.
Suggestions cannot be applied while the pull request is closed.
Suggestions cannot be applied while viewing a subset of changes.
Only one suggestion per line can be applied in a batch.
Add this suggestion to a batch that can be applied as a single commit.
Applying suggestions on deleted lines is not supported.
You must change the existing code in this line in order to create a valid suggestion.
Outdated suggestions cannot be applied.
This suggestion has been applied or marked resolved.
Suggestions cannot be applied from pending reviews.
Suggestions cannot be applied on multi-line comments.
Suggestions cannot be applied while the pull request is queued to merge.
Suggestion cannot be applied right now. Please check back later.
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Is there any risk/ambiguity with using
time.Timeas opposed to a unix msint64, for example?There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
updated to account for this risk, but used
google.protobuf.Timestampchainlink-protos changes: smartcontractkit/chainlink-protos@6acc8ef
chainlink changes: a66cc6d