Skip to content

Commit 4e76cb9

Browse files
committed
Add external_job_uri as a first-class proto field
The external job URI is already treated as special in Armada (dedicated DB column, dedicated query API endpoint), but clients have to know the magic annotation key "armadaproject.io/externalJobUri" to use it. This makes the field undiscoverable and means clients get no help from type checking or generated docs. This adds external_job_uri as a proper proto field on JobSubmitRequestItem and SubmitJob, so clients can set it directly. The server resolves the value from the proto field first and falls back to the annotation for backward compatibility. The value is mirrored into the annotation on the outgoing event so older ingesters that only read the annotation continue to work during rolling deploys. The Airflow operator is updated to use the proto field when available (guarded by hasattr for older client library versions) while continuing to set the annotation for backward compat with older servers. Signed-off-by: Dejan Zele Pejchev <pejcev.dejan@gmail.com>
1 parent 35ef598 commit 4e76cb9

14 files changed

Lines changed: 723 additions & 446 deletions

File tree

‎client/rust/src/builder.rs‎

Lines changed: 11 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -35,6 +35,7 @@ pub struct JobRequestItemBuilder<S> {
3535
labels: HashMap<String, String>,
3636
annotations: HashMap<String, String>,
3737
scheduler: String,
38+
external_job_uri: String,
3839
ingress: Vec<IngressConfig>,
3940
services: Vec<ServiceConfig>,
4041
pod_specs: Vec<PodSpec>,
@@ -56,6 +57,7 @@ impl JobRequestItemBuilder<NoPodSpec> {
5657
labels: HashMap::new(),
5758
annotations: HashMap::new(),
5859
scheduler: String::new(),
60+
external_job_uri: String::new(),
5961
ingress: Vec::new(),
6062
services: Vec::new(),
6163
pod_specs: Vec::new(),
@@ -137,6 +139,13 @@ impl<S> JobRequestItemBuilder<S> {
137139
self
138140
}
139141

142+
/// Set a URI identifying this job in an external system (e.g. Airflow).
143+
#[must_use]
144+
pub fn external_job_uri(mut self, uri: impl Into<String>) -> Self {
145+
self.external_job_uri = uri.into();
146+
self
147+
}
148+
140149
/// Replace the entire ingress config list.
141150
#[must_use]
142151
pub fn ingress(mut self, i: Vec<IngressConfig>) -> Self {
@@ -189,6 +198,7 @@ impl<S> JobRequestItemBuilder<S> {
189198
labels: self.labels,
190199
annotations: self.annotations,
191200
scheduler: self.scheduler,
201+
external_job_uri: self.external_job_uri,
192202
ingress: self.ingress,
193203
services: self.services,
194204
pod_specs: specs,
@@ -219,6 +229,7 @@ impl JobRequestItemBuilder<HasPodSpec> {
219229
ingress: self.ingress,
220230
services: self.services,
221231
scheduler: self.scheduler,
232+
external_job_uri: self.external_job_uri,
222233
// Deprecated singular fields — zeroed, not used by this builder
223234
pod_spec: None,
224235
required_node_labels: HashMap::new(),

‎client/rust/src/gen/api.rs‎

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -84,6 +84,10 @@ pub struct JobSubmitRequestItem {
8484
/// If empty, the default scheduler is used.
8585
#[prost(string, tag = "11")]
8686
pub scheduler: ::prost::alloc::string::String,
87+
/// URI identifying this job in an external system (e.g. Airflow).
88+
/// If not set, the server falls back to the "armadaproject.io/externalJobUri" annotation.
89+
#[prost(string, tag = "13")]
90+
pub external_job_uri: ::prost::alloc::string::String,
8791
}
8892
#[derive(Clone, PartialEq, ::prost::Message)]
8993
pub struct IngressConfig {

‎internal/common/constants/constants.go‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -33,6 +33,9 @@ const (
3333
FailFastAnnotation = "armadaproject.io/failFast"
3434
PoolAnnotation = "armadaproject.io/pool"
3535
ReservationTaintKey = "armadaproject.io/reservation"
36+
// ExternalJobUriAnnotation is the legacy annotation key for setting an external job URI.
37+
// Prefer the ExternalJobUri proto field on JobSubmitRequestItem / SubmitJob instead.
38+
ExternalJobUriAnnotation = "armadaproject.io/externalJobUri"
3639
)
3740

3841
var schedulingAnnotations = map[string]bool{

‎internal/lookoutingester/instructions/instructions.go‎

Lines changed: 7 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -10,6 +10,7 @@ import (
1010

1111
"github.com/armadaproject/armada/internal/common/armadacontext"
1212
"github.com/armadaproject/armada/internal/common/compress"
13+
"github.com/armadaproject/armada/internal/common/constants"
1314
"github.com/armadaproject/armada/internal/common/database/lookout"
1415
"github.com/armadaproject/armada/internal/common/eventutil"
1516
"github.com/armadaproject/armada/internal/common/ingest/metrics"
@@ -176,7 +177,12 @@ func (c *InstructionConverter) handleSubmitJob(
176177

177178
annotations := event.GetObjectMeta().GetAnnotations()
178179
userAnnotations := extractUserAnnotations(c.userAnnotationPrefix, c.blocklistAnnotations, annotations)
179-
externalJobUri := util.Truncate(annotations["armadaproject.io/externalJobUri"], maxAnnotationValLen)
180+
// Prefer the proto field; fall back to annotation for events produced before the field existed.
181+
externalJobUri := event.ExternalJobUri
182+
if externalJobUri == "" {
183+
externalJobUri = annotations[constants.ExternalJobUriAnnotation]
184+
}
185+
externalJobUri = util.Truncate(externalJobUri, maxAnnotationValLen)
180186

181187
job := model.CreateJobInstruction{
182188
JobId: event.JobId,

‎internal/lookoutingester/instructions/instructions_test.go‎

Lines changed: 57 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -14,6 +14,7 @@ import (
1414

1515
"github.com/armadaproject/armada/internal/common/armadacontext"
1616
"github.com/armadaproject/armada/internal/common/compress"
17+
"github.com/armadaproject/armada/internal/common/constants"
1718
"github.com/armadaproject/armada/internal/common/database/lookout"
1819
"github.com/armadaproject/armada/internal/common/eventutil"
1920
"github.com/armadaproject/armada/internal/common/ingest/testfixtures"
@@ -223,9 +224,9 @@ func TestConvert(t *testing.T) {
223224
}
224225
submit.GetSubmitJob().GetMainObject().GetPodSpec().GetPodSpec().PriorityClassName = priorityClass
225226
submit.GetSubmitJob().GetObjectMeta().Annotations = map[string]string{
226-
userAnnotationPrefix + "a": "0",
227-
"b": "1",
228-
"armadaproject.io/externalJobUri": "external-job-uri",
227+
userAnnotationPrefix + "a": "0",
228+
"b": "1",
229+
constants.ExternalJobUriAnnotation: "external-job-uri",
229230
}
230231
job, err := eventutil.ApiJobFromLogSubmitJob(testfixtures.UserId, []string{}, testfixtures.Queue, testfixtures.JobsetName, testfixtures.BaseTime, submit.GetSubmitJob())
231232
assert.NoError(t, err)
@@ -249,9 +250,9 @@ func TestConvert(t *testing.T) {
249250
JobProto: jobProto,
250251
PriorityClass: pointer.String(priorityClass),
251252
Annotations: map[string]string{
252-
"a": "0",
253-
"b": "1",
254-
"armadaproject.io/externalJobUri": "external-job-uri",
253+
"a": "0",
254+
"b": "1",
255+
constants.ExternalJobUriAnnotation: "external-job-uri",
255256
},
256257
ExternalJobUri: "external-job-uri",
257258
}
@@ -538,14 +539,63 @@ func TestConvert(t *testing.T) {
538539
}
539540
}
540541

542+
func TestExternalJobUriProtoFieldPreferred(t *testing.T) {
543+
tests := map[string]struct {
544+
protoField string
545+
annotations map[string]string
546+
expected string
547+
}{
548+
"Proto field takes priority over annotation": {
549+
protoField: "airflow://dag/task/run/0",
550+
annotations: map[string]string{
551+
constants.ExternalJobUriAnnotation: "old-value",
552+
},
553+
expected: "airflow://dag/task/run/0",
554+
},
555+
"Falls back to annotation when proto field empty": {
556+
protoField: "",
557+
annotations: map[string]string{
558+
constants.ExternalJobUriAnnotation: "airflow://dag/task/run/0",
559+
},
560+
expected: "airflow://dag/task/run/0",
561+
},
562+
"Empty when neither set": {
563+
protoField: "",
564+
annotations: map[string]string{},
565+
expected: "",
566+
},
567+
}
568+
569+
for name, tc := range tests {
570+
t.Run(name, func(t *testing.T) {
571+
submit, err := testfixtures.DeepCopy(testfixtures.Submit)
572+
require.NoError(t, err)
573+
submit.GetSubmitJob().ExternalJobUri = tc.protoField
574+
submit.GetSubmitJob().GetObjectMeta().Annotations = tc.annotations
575+
576+
events := &utils.EventsWithIds[*armadaevents.EventSequence]{
577+
Events: []*armadaevents.EventSequence{
578+
testfixtures.NewEventSequence(submit),
579+
},
580+
MessageIds: []pulsar.MessageID{pulsarutils.NewMessageId(1)},
581+
}
582+
583+
converter := NewInstructionConverter(metrics.Get().Metrics, userAnnotationPrefix, []string{}, &compress.NoOpCompressor{})
584+
instructionSet := converter.Convert(armadacontext.TODO(), events)
585+
require.Len(t, instructionSet.JobsToCreate, 1)
586+
assert.Equal(t, tc.expected, instructionSet.JobsToCreate[0].ExternalJobUri)
587+
})
588+
}
589+
}
590+
541591
func TestTruncatesStringsThatAreTooLong(t *testing.T) {
542592
longString := strings.Repeat("x", 4000)
543593

544594
submit, err := testfixtures.DeepCopy(testfixtures.Submit)
545595
assert.NoError(t, err)
546596
submit.GetSubmitJob().GetMainObject().GetPodSpec().GetPodSpec().PriorityClassName = longString
547597
submit.GetSubmitJob().GetObjectMeta().Annotations = map[string]string{
548-
"armadaproject.io/externalJobUri": longString,
598+
constants.ExternalJobUriAnnotation: longString,
549599
}
550600

551601
leased, err := testfixtures.DeepCopy(testfixtures.Leased)

‎internal/server/submit/conversion/conversions.go‎

Lines changed: 19 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -8,6 +8,7 @@ import (
88
networking "k8s.io/api/networking/v1"
99

1010
"github.com/armadaproject/armada/internal/common"
11+
"github.com/armadaproject/armada/internal/common/constants"
1112
log "github.com/armadaproject/armada/internal/common/logging"
1213
armadaslices "github.com/armadaproject/armada/internal/common/slices"
1314
"github.com/armadaproject/armada/internal/common/util"
@@ -29,13 +30,27 @@ func SubmitJobFromApiRequest(
2930
priority := PriorityAsInt32(jobReq.GetPriority())
3031
ingressesAndServices := convertIngressesAndServices(config, jobReq, jobId, jobSetId, queue, owner)
3132

33+
// Resolve externalJobUri: prefer the proto field, fall back to annotation.
34+
annotations := jobReq.GetAnnotations()
35+
externalJobUri := jobReq.GetExternalJobUri()
36+
if externalJobUri == "" {
37+
externalJobUri = annotations[constants.ExternalJobUriAnnotation]
38+
}
39+
// Mirror into annotation so old ingesters (that only read the annotation) still work during rolling deploys.
40+
if externalJobUri != "" {
41+
if annotations == nil {
42+
annotations = make(map[string]string)
43+
}
44+
annotations[constants.ExternalJobUriAnnotation] = externalJobUri
45+
}
46+
3247
msg := &armadaevents.SubmitJob{
3348
JobId: jobId,
3449
DeduplicationId: jobReq.GetClientId(),
3550
Priority: priority,
3651
ObjectMeta: &armadaevents.ObjectMeta{
3752
Namespace: jobReq.GetNamespace(),
38-
Annotations: jobReq.GetAnnotations(),
53+
Annotations: annotations,
3954
Labels: jobReq.GetLabels(),
4055
},
4156
MainObject: &armadaevents.KubernetesMainObject{
@@ -45,8 +60,9 @@ func SubmitJobFromApiRequest(
4560
},
4661
},
4762
},
48-
Objects: ingressesAndServices,
49-
Scheduler: jobReq.Scheduler,
63+
Objects: ingressesAndServices,
64+
Scheduler: jobReq.Scheduler,
65+
ExternalJobUri: externalJobUri,
5066
}
5167

5268
postProcess(msg, config)

‎internal/server/submit/conversion/conversions_test.go‎

Lines changed: 58 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -5,10 +5,12 @@ import (
55
"testing"
66

77
"github.com/stretchr/testify/assert"
8+
"github.com/stretchr/testify/require"
89
v1 "k8s.io/api/core/v1"
910
networking "k8s.io/api/networking/v1"
1011
"k8s.io/apimachinery/pkg/api/resource"
1112

13+
"github.com/armadaproject/armada/internal/common/constants"
1214
"github.com/armadaproject/armada/internal/server/configuration"
1315
"github.com/armadaproject/armada/internal/server/submit/testfixtures"
1416
"github.com/armadaproject/armada/pkg/api"
@@ -907,6 +909,62 @@ func jobSubmitRequestItemWithClassicInitContainer() *api.JobSubmitRequestItem {
907909
return req
908910
}
909911

912+
func TestExternalJobUri(t *testing.T) {
913+
tests := map[string]struct {
914+
protoField string
915+
annotations map[string]string
916+
expectedUri string
917+
expectMirrored bool
918+
}{
919+
"Proto field takes priority over annotation": {
920+
protoField: "airflow://dag/task/run/0",
921+
annotations: map[string]string{constants.ExternalJobUriAnnotation: "old-annotation-value"},
922+
expectedUri: "airflow://dag/task/run/0",
923+
expectMirrored: true,
924+
},
925+
"Falls back to annotation when proto field empty": {
926+
protoField: "",
927+
annotations: map[string]string{constants.ExternalJobUriAnnotation: "airflow://dag/task/run/0"},
928+
expectedUri: "airflow://dag/task/run/0",
929+
expectMirrored: true,
930+
},
931+
"Empty when neither set": {
932+
protoField: "",
933+
annotations: map[string]string{},
934+
expectedUri: "",
935+
expectMirrored: false,
936+
},
937+
"Proto field mirrored into annotation": {
938+
protoField: "airflow://dag/task/run/0",
939+
annotations: nil,
940+
expectedUri: "airflow://dag/task/run/0",
941+
expectMirrored: true,
942+
},
943+
}
944+
945+
for name, tc := range tests {
946+
t.Run(name, func(t *testing.T) {
947+
jobReq := testfixtures.JobSubmitRequestItem(1)
948+
jobReq.Annotations = tc.annotations
949+
jobReq.ExternalJobUri = tc.protoField
950+
951+
msg := SubmitJobFromApiRequest(
952+
jobReq,
953+
testfixtures.DefaultSubmissionConfig(),
954+
testfixtures.DefaultJobset,
955+
testfixtures.DefaultQueue.Name,
956+
testfixtures.DefaultOwner,
957+
testfixtures.TestUlidGenerator(),
958+
)
959+
960+
require.Equal(t, tc.expectedUri, msg.ExternalJobUri)
961+
if tc.expectMirrored {
962+
assert.Equal(t, tc.expectedUri, msg.GetObjectMeta().GetAnnotations()[constants.ExternalJobUriAnnotation])
963+
}
964+
})
965+
}
966+
}
967+
910968
func SubmitJobMsgWithK8sObjects(objects []*armadaevents.KubernetesObject, initContainers ...v1.Container) *armadaevents.SubmitJob {
911969
submitMsg := testfixtures.SubmitJob(1)
912970
submitMsg.Objects = objects

‎pkg/api/api.swagger.go‎

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1881,6 +1881,10 @@ func SwaggerJsonTemplate() string {
18811881
" \"clientId\": {\n" +
18821882
" \"type\": \"string\"\n" +
18831883
" },\n" +
1884+
" \"externalJobUri\": {\n" +
1885+
" \"description\": \"URI identifying this job in an external system (e.g. Airflow).\\nIf not set, the server falls back to the \\\"armadaproject.io/externalJobUri\\\" annotation.\",\n" +
1886+
" \"type\": \"string\"\n" +
1887+
" },\n" +
18841888
" \"ingress\": {\n" +
18851889
" \"type\": \"array\",\n" +
18861890
" \"items\": {\n" +

‎pkg/api/api.swagger.json‎

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1870,6 +1870,10 @@
18701870
"clientId": {
18711871
"type": "string"
18721872
},
1873+
"externalJobUri": {
1874+
"description": "URI identifying this job in an external system (e.g. Airflow).\nIf not set, the server falls back to the \"armadaproject.io/externalJobUri\" annotation.",
1875+
"type": "string"
1876+
},
18731877
"ingress": {
18741878
"type": "array",
18751879
"items": {

0 commit comments

Comments
 (0)