Skip to content

Commit eca4986

Browse files
additional retry for stuff that is seeing transient errors/flakiness
1 parent eaa64a1 commit eca4986

3 files changed

Lines changed: 56 additions & 11 deletions

File tree

tests/integration/metadata/test_schema_service.py

Lines changed: 16 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,7 @@
55
from conductor.client.http.api.schema_resource_api import SchemaResourceApi
66
from conductor.client.http.models.schema_def import SchemaDef, SchemaType
77
from conductor.client.orkes.orkes_schema_client import OrkesSchemaClient
8+
from tests.integration.retry_helpers import retry_on_transient, retry_on_status
89

910
SCHEMA_NAME = 'ut_schema'
1011
SCHEMA_VERSION = 1
@@ -41,33 +42,40 @@ def test_init(self):
4142
self.assertIsInstance(self.schema_client.schemaApi, SchemaResourceApi, message)
4243

4344
def test_registerSchema(self):
44-
self.schema_client.register_schema(self.schemaDef)
45-
response = self.schema_client.schemaApi.get_schema_by_name_and_version(name=SCHEMA_NAME, version=SCHEMA_VERSION)
45+
retry_on_transient(self.schema_client.register_schema, self.schemaDef)
46+
# A GET right after register can briefly 404 until the write propagates
47+
# on the shared dev server; retry the read rather than fail the test.
48+
response = retry_on_status(
49+
self.schema_client.schemaApi.get_schema_by_name_and_version,
50+
name=SCHEMA_NAME, version=SCHEMA_VERSION)
4651
self.assertEqual(response.name, SCHEMA_NAME)
4752
self.assertEqual(response.version, SCHEMA_VERSION)
4853
self.assertEqual(response.type, SchemaType.JSON)
4954

5055
def test_getSchema(self):
51-
self.schema_client.register_schema(self.schemaDef)
52-
schema = self.schema_client.get_schema(SCHEMA_NAME, SCHEMA_VERSION)
56+
retry_on_transient(self.schema_client.register_schema, self.schemaDef)
57+
# A GET right after register can briefly 404 until the write propagates
58+
# on the shared dev server; retry the read rather than fail the test.
59+
schema = retry_on_status(self.schema_client.get_schema,
60+
SCHEMA_NAME, SCHEMA_VERSION)
5361
self.assertEqual(schema.name, SCHEMA_NAME)
5462
self.assertEqual(schema.version, SCHEMA_VERSION)
5563

5664
def test_getAllSchemas(self):
5765
schemaDef2 = SchemaDef(name='ut_schema_2', version=1, type=SchemaType.JSON, data=schema, external_ref='http://example.com/2')
58-
self.schema_client.register_schema(self.schemaDef)
59-
self.schema_client.register_schema(schemaDef2)
66+
retry_on_transient(self.schema_client.register_schema, self.schemaDef)
67+
retry_on_transient(self.schema_client.register_schema, schemaDef2)
6068
schemas = self.schema_client.get_all_schemas()
6169
self.assertGreaterEqual(len(schemas), 2)
6270

6371
def test_deleteSchema(self):
64-
self.schema_client.register_schema(self.schemaDef)
72+
retry_on_transient(self.schema_client.register_schema, self.schemaDef)
6573
self.schema_client.delete_schema(SCHEMA_NAME, SCHEMA_VERSION)
6674
with self.assertRaises(Exception):
6775
self.schema_client.get_schema(SCHEMA_NAME, SCHEMA_VERSION)
6876

6977
def test_deleteSchemaByName(self):
70-
self.schema_client.register_schema(self.schemaDef)
78+
retry_on_transient(self.schema_client.register_schema, self.schemaDef)
7179
self.schema_client.delete_schema_by_name(SCHEMA_NAME)
7280
with self.assertRaises(Exception):
7381
self.schema_client.get_schema(SCHEMA_NAME, SCHEMA_VERSION)

tests/integration/retry_helpers.py

Lines changed: 28 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -93,6 +93,34 @@ def retry_on_transient(func, *args, retries=4, base_delay=1, **kwargs):
9393
raise
9494

9595

96+
def retry_on_status(func, *args, statuses=(404,), retries=5, base_delay=1.0,
97+
max_delay=None, **kwargs):
98+
"""Retry ``func(*args, **kwargs)`` on an ``ApiException`` whose status is in
99+
``statuses`` (in addition to transient status-0 blips), with capped
100+
exponential backoff. Any other error raises immediately.
101+
102+
Intended for read-after-write races against the shared dev server: a GET
103+
issued right after a register/update can briefly 404 until the write
104+
propagates. This is *per-request* retry, so use it only for idempotent reads
105+
(or writes that are safe to repeat).
106+
"""
107+
for attempt in range(retries):
108+
try:
109+
return func(*args, **kwargs)
110+
except ApiException as e:
111+
retryable = is_transient(e) or e.status in statuses
112+
if retryable and attempt < retries - 1:
113+
delay = base_delay * (2 ** attempt)
114+
if max_delay is not None:
115+
delay = min(delay, max_delay)
116+
logger.warning(
117+
'retryable (%s) API error (attempt %d/%d): %s; retrying in '
118+
'%.1fs', e.status, attempt + 1, retries, e, delay)
119+
time.sleep(delay)
120+
continue
121+
raise
122+
123+
96124
def retry_scenario(label, func, *args, deadline=None,
97125
base_delay=DEFAULT_BASE_DELAY_SECONDS,
98126
max_delay=DEFAULT_MAX_DELAY_SECONDS, **kwargs):

tests/integration/test_v2_fallback_intg.py

Lines changed: 12 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -25,6 +25,7 @@
2525
from conductor.client.orkes.orkes_metadata_client import OrkesMetadataClient
2626
from conductor.client.orkes.orkes_workflow_client import OrkesWorkflowClient
2727
from conductor.client.worker.worker_task import worker_task
28+
from tests.integration.retry_helpers import retry_on_transient
2829

2930
logger = logging.getLogger(__name__)
3031

@@ -84,10 +85,14 @@ def test_0_register_workflow(self):
8485

8586
workflow = WorkflowDef(name=WORKFLOW_NAME, version=WORKFLOW_VERSION)
8687
workflow._tasks = tasks
88+
# Retry registration on a transient (status 0) transport blip against the
89+
# shared dev server so a dropped connection doesn't fail the suite.
8790
try:
88-
self.metadata_client.update_workflow_def(workflow, overwrite=True)
91+
retry_on_transient(self.metadata_client.update_workflow_def,
92+
workflow, overwrite=True)
8993
except Exception:
90-
self.metadata_client.register_workflow_def(workflow, overwrite=True)
94+
retry_on_transient(self.metadata_client.register_workflow_def,
95+
workflow, overwrite=True)
9196
print(f"\n Registered workflow '{WORKFLOW_NAME}' with {len(tasks)} tasks")
9297

9398
def test_1_workflows_complete_with_v2_or_fallback(self):
@@ -128,7 +133,11 @@ def _run_workers():
128133
req.name = WORKFLOW_NAME
129134
req.version = WORKFLOW_VERSION
130135
req.input = {"run_index": i}
131-
wf_id = self.workflow_client.start_workflow(start_workflow_request=req)
136+
# A status-0 blip here means no response arrived, so no id was
137+
# returned; retrying gets a fresh attempt (any orphaned run just
138+
# completes untracked and doesn't affect the tracked-id count).
139+
wf_id = retry_on_transient(self.workflow_client.start_workflow,
140+
start_workflow_request=req)
132141
workflow_ids.append(wf_id)
133142

134143
print(f"\n Submitted {len(workflow_ids)} workflows")

0 commit comments

Comments
 (0)