Repository navigation
feat(OMN-1934): Add onex-infra-test integration test CLI - #261
Conversation
Scaffold the `onex-infra-test` CLI for end-to-end validation of the ONEX registration pipeline. Implements all 8 tasks from Phase 5: - P5.1: CLI scaffold using Click with Rich output - P5.2: `env up/down` wrapping docker-compose.e2e.yml with health wait - P5.3: `verify registry` checks Consul + PostgreSQL state - P5.4: `verify topics` validates ONEX 5-segment naming compliance - P5.5: `verify snapshots` checks compacted snapshot topic - P5.6: `verify idempotency` publishes N events, asserts <= 1 record - P5.7: `run --suite smoke` orchestrates happy-path verification - P5.8: `run --suite failure` tests runtime kill/restart recovery Also adds `introspect` command to publish test introspection events via rpk, and `run --suite idempotency` for duplicate detection. Entry point: `onex-infra-test` (registered in pyproject.toml)
- Remove unused imports (sys, asyncio) - Log exception details in bare except blocks (Consul, PostgreSQL) - Include exception type in all error messages for consistency - Document hardcoded test defaults as E2E-only
📝 WalkthroughWalkthroughAdds a new CLI package Changes
Sequence Diagram(s)sequenceDiagram
participant User as User/CLI
participant CLI as onex-infra-test CLI
participant Compose as Docker Compose
participant Kafka as Kafka (rpk)
participant DB as PostgreSQL
participant Consul as Consul
User->>CLI: invoke command (env / introspect / verify / run)
CLI->>Compose: start/stop services (compose-file, project-name)
Compose-->>CLI: status / health
CLI->>Kafka: publish introspection (rpk -> topic)
Kafka-->>CLI: publish result
CLI->>DB: poll registration_projections for node_id
DB-->>CLI: registration found / not found
CLI->>Consul: query KV / services
Consul-->>CLI: registry entries
CLI->>User: display PASS / WARN / FAIL
Estimated code review effort🎯 4 (Complex) | ⏱️ ~45 minutes Possibly related PRs
Poem
🚥 Pre-merge checks | ✅ 3✅ Passed checks (3 passed)
✏️ Tip: You can configure your own custom pre-merge checks in the settings. ✨ Finishing touches
🧪 Generate unit tests (beta)
Comment |
There was a problem hiding this comment.
Actionable comments posted: 6
🤖 Fix all issues with AI agents
In `@src/omnibase_infra/cli/infra_test/introspect.py`:
- Around line 27-58: The function _build_introspection_payload currently ignores
the version parameter and hardcodes node_version; update
_build_introspection_payload to parse the version string argument (e.g.,
"1.2.3") into integers for major/minor/patch and populate "node_version" with
those values, handling missing segments by treating them as 0 and raising or
validating on non-numeric input; reference the function name
_build_introspection_payload and the "node_version" dict to locate where to
replace the hardcoded {"major":1,"minor":0,"patch":0} with the parsed values
from the version parameter.
In `@src/omnibase_infra/cli/infra_test/run_suite.py`:
- Around line 446-449: The inline message "Waiting for runtime health (30s)..."
in run_suite.py is inconsistent with the actual wait call; update the code so
the console.print and time.sleep match: either change time.sleep(10) to
time.sleep(30) to match the 30s message, or change the printed message to
"(10s)" to reflect the current time.sleep(10) behavior; modify the console.print
and/or time.sleep calls (the console.print(...) and time.sleep(...) lines)
accordingly so they are consistent.
- Around line 105-127: The polling loop may leak DB connections if an exception
occurs after psycopg2.connect() but before conn.close(); update the block that
calls psycopg2.connect(dsn, ...) and then cur = conn.cursor() / cur.execute(...)
to ensure cleanup by using context managers (with psycopg2.connect(...) as conn:
and with conn.cursor() as cur:) or a try/finally that always calls cur.close()
and conn.close(); specifically modify the connect/cursor/execute section that
uses dsn, node_id, timeout, console so that cur.close() and conn.close() are
guaranteed on error and successful paths.
- Around line 309-323: The current try/except around psycopg2.connect leaks
resources; update the block in run_suite to use context managers: open the
connection with "with psycopg2.connect(dsn, connect_timeout=5) as conn:" and the
cursor with "with conn.cursor() as cur:" around the cur.execute/fetchone logic
so conn and cur are closed automatically, preserve the existing exception
handling that logs the error (console.print with type(e).__name__ and e) and
then raises SystemExit(1).
- Around line 461-484: Replace the manual connect/cursor open/close in the
post-restart check with context managers to avoid leaking connections: use
psycopg2.connect(...) as conn and conn.cursor() as cur to run the SELECT against
registration_projections for entity_id node_a, fetchone and let the context
managers close the cursor and connection automatically; keep the existing
SystemExit and general Exception handling but remove explicit
cur.close()/conn.close() calls and ensure the same console prints and exit
behavior remain.
In `@tests/unit/cli/infra_test/test_introspect.py`:
- Around line 46-52: The test reveals _build_introspection_payload currently
ignores its version parameter and always sets node_version to
{"major":1,"minor":0,"patch":0}; update _build_introspection_payload to respect
the version argument (e.g., accept a string like "2.3.1" or a tuple/dict) by
parsing it and populating node_version accordingly (major/minor/patch ints) and
fallback to the default 1.0.0 when version is None or invalid, and then add a
unit test (besides test_has_version_semver) that passes a custom version (e.g.,
"2.3.1") and asserts node_version reflects 2, 3, 1; alternatively, if you prefer
not to support custom versions, remove the unused version parameter from
_build_introspection_payload and adjust tests to match the simplified signature.
🧹 Nitpick comments (10)
src/omnibase_infra/cli/infra_test/introspect.py (2)
113-128: Add a timeout to prevent indefinite hangs.The
subprocess.runcall lacks a timeout. Ifrpkhangs (e.g., broker unreachable without proper timeout handling in rpk itself), the CLI will block indefinitely.⏱️ Proposed fix to add timeout
result = subprocess.run( [ "rpk", "topic", "produce", topic, "--brokers", resolved_broker, "-k", str(nid), ], input=payload_json, capture_output=True, text=True, check=False, + timeout=30, )
96-98: Redundant environment variable resolution.The
--brokeroption already specifiesenvvar="KAFKA_BOOTSTRAP_SERVERS"(line 81), so Click will automatically resolve from the environment. The manualos.getenv("KAFKA_BOOTSTRAP_SERVERS")check is redundant.♻️ Simplified broker resolution
- resolved_broker: str = ( - broker or os.getenv("KAFKA_BOOTSTRAP_SERVERS") or "localhost:29092" - ) + resolved_broker: str = broker or "localhost:29092"src/omnibase_infra/cli/infra_test/verify.py (4)
29-31: Duplicate helper function.
_get_broker()is duplicated identically inrun_suite.py(lines 23-25). Consider extracting to a shared utility module to avoid drift.
96-101: Add timeout tosubprocess.runcalls.Multiple
subprocess.runcalls in this file (lines 96-101, 172-177, 192-209, 260-275) lack timeouts. Ifrpkhangs, the CLI will block indefinitely. Consider addingtimeout=30consistently.
385-408: Use context manager for database connections.The psycopg2 connection is not closed if an exception occurs between
connect()and the explicitclose()calls. Using a context manager ensures proper cleanup.🔒 Proposed fix using context manager
try: import psycopg2 - conn = psycopg2.connect(dsn, connect_timeout=5) - cur = conn.cursor() - - if node_id: - cur.execute( - "SELECT entity_id, current_state, node_type, updated_at " - "FROM registration_projections WHERE entity_id = %s", - (node_id,), - ) - else: - cur.execute( - "SELECT entity_id, current_state, node_type, updated_at " - "FROM registration_projections ORDER BY updated_at DESC LIMIT 20" - ) - - rows = cur.fetchall() - cur.close() - conn.close() + with psycopg2.connect(dsn, connect_timeout=5) as conn: + with conn.cursor() as cur: + if node_id: + cur.execute( + "SELECT entity_id, current_state, node_type, updated_at " + "FROM registration_projections WHERE entity_id = %s", + (node_id,), + ) + else: + cur.execute( + "SELECT entity_id, current_state, node_type, updated_at " + "FROM registration_projections ORDER BY updated_at DESC LIMIT 20" + ) + rows = cur.fetchall() except Exception as e:The same pattern applies to
_count_postgres_registrations(lines 437-450).
220-304: Hardcoded 5-second sleep may cause flaky tests.Line 285 uses a fixed
time.sleep(5)to wait for registration processing. This could be too short in slow CI environments or wastefully long in fast environments. Consider making this configurable or implementing a polling approach with exponential backoff.tests/unit/cli/infra_test/test_verify_topics.py (1)
17-18: Missing@pytest.mark.unitmarker.Per coding guidelines, test files should use pytest markers for classification. Add the unit marker to the test class.
🏷️ Proposed fix to add pytest marker
+@pytest.mark.unit class TestOnexTopicPattern: """Test ONEX 5-segment topic naming regex."""As per coding guidelines: "Use pytest markers:
@pytest.mark.unit,@pytest.mark.integration,@pytest.mark.slow,@pytest.mark.chaos,@pytest.mark.serialfor test classification"tests/unit/cli/infra_test/test_introspect.py (1)
41-44: Weak assertion on timestamp—only checks truthiness.The test asserts
payload["timestamp"]is truthy but doesn't verify it's a valid ISO 8601 string. Consider validating the format to catch malformed timestamps.💡 Optional: Validate ISO 8601 format
+ from datetime import datetime + def test_has_timestamp(self) -> None: """Payload includes a timestamp.""" payload = _build_introspection_payload() - assert payload["timestamp"] + ts = payload["timestamp"] + assert ts + # Validate ISO 8601 format + datetime.fromisoformat(ts.replace("Z", "+00:00"))src/omnibase_infra/cli/infra_test/run_suite.py (1)
196-218: Consider extracting duplicated environment verification logic.The environment verification check (docker compose ps) is repeated nearly identically in all three suite functions. A shared helper would reduce duplication and ensure consistent behavior.
♻️ Suggested helper extraction
def _verify_environment_running( compose_file: str, project_name: str, require_runtime: bool = False, ) -> None: """Verify E2E environment is running, exit if not.""" result = subprocess.run( [ "docker", "compose", "-f", compose_file, "-p", project_name, "ps", "--status" if not require_runtime else "--services", "running" if not require_runtime else None, ], capture_output=True, text=True, check=False, ) # Handle validation logic...Also applies to: 276-295, 358-380
tests/unit/cli/infra_test/test_cli_scaffold.py (1)
16-122: Missing@pytest.mark.unitdecorator on test class.Consistent with the earlier review comment on
test_introspect.py, this test class should include the@pytest.mark.unitmarker unless auto-applied via conftest.As per coding guidelines: "Use pytest markers:
@pytest.mark.unit,@pytest.mark.integration,@pytest.mark.slow,@pytest.mark.chaos,@pytest.mark.serialfor test classification".
Extract duplicated _get_broker(), _get_consul_addr(), and _get_postgres_dsn() into a shared _helpers.py module. Add str() casts and isinstance() guards in test_introspect.py to satisfy mypy with the dict[str, object] return type. Remove redundant os.getenv broker resolution in introspect command (Click already handles the envvar).
…yload
The version parameter was accepted but ignored — the function always
hardcoded node_version to {1, 0, 0}. Now parses the version string
into major/minor/patch components.
Review iteration: 1/10
- Fix connection leaks: use try/finally with conn.close() for all psycopg2 connections in run_suite.py and verify.py - Add subprocess timeout to all subprocess.run calls to prevent CI hangs - Extract duplicated env verification into _verify_env_running() helper - Fix comment/code mismatch (30s comment vs 10s sleep) in failure suite - Centralize broker fallback via get_broker() in introspect.py - Add @pytest.mark.unit to all 3 test classes - Add test_custom_version to verify version parameter affects output - Strengthen timestamp assertion to validate ISO 8601 format
- Fixed: introspect.py:48 - Add try/except for version parsing with click.BadParameter on non-numeric input - Fixed: _helpers.py:50 - URL-encode user/password in get_postgres_dsn() via urllib.parse.quote_plus - Added: test_helpers.py - Unit tests for get_broker, get_consul_addr, get_postgres_dsn (including special character encoding) - Added: test_introspect.py - Test for invalid version rejection Review iteration: 1/10
There was a problem hiding this comment.
Actionable comments posted: 3
🤖 Fix all issues with AI agents
In `@src/omnibase_infra/cli/infra_test/_helpers.py`:
- Around line 11-12: The DSN construction in this module builds a Postgres URI
using raw username/password which fails for special characters; import
urllib.parse and URL-encode the credentials with urllib.parse.quote_plus when
constructing the DSN (apply to the DSN builder used around lines 35-50). Find
the function or variable that assembles the Postgres URI (e.g.,
get_postgres_dsn, build_dsn, or the DSN variable) and replace raw username and
password insertion with quote_plus(username) and quote_plus(password), ensuring
the final URI uses the encoded values.
In `@src/omnibase_infra/cli/infra_test/introspect.py`:
- Around line 47-50: The version parsing currently does int(parts[i]) directly
which raises ValueError for non-numeric inputs; update the parsing of
version/parts to validate each part (e.g., str.isdigit() or try/except around
int conversion) and handle conversion failures by falling back to the default
for that segment (major=1, minor=0, patch=0) and emitting a friendly CLI error
(raise SystemExit or use the CLI framework's bad-parameter error) when the
entire --version is invalid; apply the same guarded logic to the other
occurrence referenced at lines 106-107 (the variables: version, parts, major,
minor, patch).
In `@src/omnibase_infra/cli/infra_test/verify.py`:
- Around line 292-413: The checks in _verify_consul_registry and
_verify_postgres_registry currently treat an overall empty result as non-failure
but do not fail when a specific node_id is requested but no matching entries
exist; update both functions to explicitly fail when node_id is provided and no
entries match the filter. In _verify_consul_registry, after building
onex_services or keys and applying the node_id filter (the loop that uses "if
node_id and node_id not in ...: continue"), detect if no items were
added/matched and console.print a red error mentioning the missing node_id and
return False; also adjust the printed table/count to reflect only matched items.
In _verify_postgres_registry, after fetching rows from registration_projections,
filter rows by node_id (or use the node-specific query already present) and if
node_id is provided but the filtered rows list is empty, console.print a red
error and return False; update the final count to use the number of matched rows
when printing. Ensure you reference the functions _verify_consul_registry and
_verify_postgres_registry and the variables onex_services, keys, and rows when
making these changes.
…dation - Fixed: verify.py - verify registry --node-id returns FAIL when specified node not found (Consul KV, services, and PostgreSQL paths) - Fixed: _helpers.py - Validate POSTGRES_HOST has no '@', POSTGRES_PORT is numeric, and CONSUL_SCHEME is http/https - Fixed: run_suite.py - Narrow except Exception to psycopg2.Error in _wait_for_registration polling loop - Added: test_helpers.py - Tests for invalid scheme, host, and port Review iteration: 2/10
- Fixed: verify.py - Replace broad except Exception with httpx.HTTPError for Consul calls and psycopg2.Error for PostgreSQL calls - Fixed: run_suite.py - Move psycopg2 imports before try blocks in idempotency and failure suites; narrow except to psycopg2.Error - Fixed: _helpers.py - URL-encode POSTGRES_DATABASE in DSN via quote_plus for consistency with user/password encoding Review iteration: 3/10
…nses - Fixed: verify.py - Widen Consul exception handlers to catch (httpx.HTTPError, ValueError) so json.JSONDecodeError (a ValueError subclass) from non-JSON 2xx responses is handled gracefully Review iteration: 4/10
There was a problem hiding this comment.
Actionable comments posted: 1
🤖 Fix all issues with AI agents
In `@src/omnibase_infra/cli/infra_test/run_suite.py`:
- Around line 305-313: The idempotency check currently allows count==0 to pass;
update the Idempotency suite in run_suite.py to first poll for at least one
record (reuse _wait_for_registration(dsn, node_id, timeout=15) or equivalent)
and fail early if no registration appears, then fetch the count and assert there
is exactly one record (or assert count > 0 and count <= 1) instead of allowing
zero; ensure you raise SystemExit(1) with a clear message when the wait times
out so the suite doesn't falsely pass when no events were processed.
🧹 Nitpick comments (1)
src/omnibase_infra/cli/infra_test/verify.py (1)
376-377: Minor: Count message shows total keys instead of matched count.When filtering by
--node-id, line 376 displayslen(keys)(total keys) rather thanmatched(filtered count). This could be confusing since only matched keys appear in the table.💡 Suggested fix
- console.print(f" [green]{len(keys)} service key(s) found.[/green]") + display_count = matched if node_id else len(keys) + console.print(f" [green]{display_count} service key(s) found.[/green]")
| console.print(f" Records found: {count}") | ||
|
|
||
| if count <= 1: | ||
| console.print("[bold green]Idempotency suite: PASS[/bold green]") | ||
| else: | ||
| console.print( | ||
| f"[bold red]Idempotency suite: FAIL (expected <= 1, found {count})[/bold red]" | ||
| ) | ||
| raise SystemExit(1) |
There was a problem hiding this comment.
Idempotency test may pass falsely when no events are processed.
The test asserts count <= 1, which passes when count == 0. If events weren't processed within the 5-second settle time (e.g., consumer lag, service not running), the test would pass without verifying idempotency. Consider polling until at least one record exists (similar to _wait_for_registration in the smoke suite), then asserting count == 1.
🛠️ Suggested approach
console.print(f" Records found: {count}")
- if count <= 1:
+ if count == 1:
console.print("[bold green]Idempotency suite: PASS[/bold green]")
+ elif count == 0:
+ console.print(
+ "[bold red]Idempotency suite: FAIL "
+ "(no records found - events may not have been processed)[/bold red]"
+ )
+ raise SystemExit(1)
else:
console.print(
f"[bold red]Idempotency suite: FAIL (expected <= 1, found {count})[/bold red]"
)
raise SystemExit(1)Alternatively, poll for at least one record before checking the count:
# Wait for at least one record to appear
if not _wait_for_registration(dsn, node_id, timeout=15):
console.print("[bold red]No registration found within timeout.[/bold red]")
raise SystemExit(1)
# Then check for duplicates
count = ...🤖 Prompt for AI Agents
In `@src/omnibase_infra/cli/infra_test/run_suite.py` around lines 305 - 313, The
idempotency check currently allows count==0 to pass; update the Idempotency
suite in run_suite.py to first poll for at least one record (reuse
_wait_for_registration(dsn, node_id, timeout=15) or equivalent) and fail early
if no registration appears, then fetch the count and assert there is exactly one
record (or assert count > 0 and count <= 1) instead of allowing zero; ensure you
raise SystemExit(1) with a clear message when the wait times out so the suite
doesn't falsely pass when no events were processed.
…checks - Convert all psycopg2.connect() to context managers in run_suite.py (3 sites) and verify.py (2 sites) to prevent connection leaks on exception - Add version input guards in introspect.py: reject empty strings, >3 segments, and negative numbers - Strengthen timestamp display in verify.py with type/range validation instead of bare truthiness check - Add 5 edge-case tests for version parsing (partial, extra segments, negative, empty)
There was a problem hiding this comment.
Actionable comments posted: 2
🤖 Fix all issues with AI agents
In `@src/omnibase_infra/cli/infra_test/run_suite.py`:
- Around line 271-277: The idempotency test is generating new payloads each
iteration so events differ; instead construct a single introspection payload
once and reuse it for all publishes. Update the loop in run_suite.py to build
the payload before the for i in range(repetitions) loop and call
_publish_introspection(broker, topic, node_id, payload) (or modify
_publish_introspection to accept an optional payload) so
correlation_id/timestamp remain identical across all iterations; keep the
existing failure handling and console prints unchanged.
In `@src/omnibase_infra/cli/infra_test/verify.py`:
- Around line 268-285: The idempotency check currently treats count == 0 as a
pass; update the logic in verify.py where get_postgres_dsn() and
_count_postgres_registrations(dsn, node_id) are used so that count == None still
raises an error, count == 1 is the only success case, and both count == 0 and
count > 1 are treated as failures: for count == 0 print a clear "[bold red]No
registration records observed (possible processing failure or consumer
lag).[/bold red]" message and exit with SystemExit(1), and for count > 1 keep
the existing duplicate-record failure path.
…sitives - Fix psycopg2 connection leaks with explicit conn.close() in try/finally - Ensure idempotency suite publishes byte-identical payloads - Fail when count=0 in idempotency checks (no false positives) - Show matched vs total count when filtering Consul KV by node-id
…ous-truth guard - Fixed: run_suite.py:308,469 - Wrap psycopg2.connect() inside try/except so connection failures produce clean CLI errors instead of raw tracebacks - Fixed: _helpers.py:31 - Add port.isdigit() guard to get_consul_addr(), matching the existing validation in get_postgres_dsn() - Fixed: env.py:108 - Guard _wait_for_healthy against vacuous-truth when no service lines contain tab separators (same bug class as e58a17b) Review iteration: 1/10
There was a problem hiding this comment.
Actionable comments posted: 1
🤖 Fix all issues with AI agents
In `@src/omnibase_infra/cli/infra_test/verify.py`:
- Around line 433-455: The closing Rich markup tag for state styling is
incorrectly built using state_style.strip('[') which yields "green]" and creates
invalid markup; update the construction in the table.add_row call that builds
the state cell (the variables state_style and state are used there) so you strip
both brackets (e.g., strip('[]')) or otherwise extract the bare style name from
state_style and then concatenate it into a proper closing tag like
f"[/{style_name}]" so the final string is a valid closing tag for the opening
state_style.
| for row in rows: | ||
| entity_id, state, ntype, updated = row | ||
| state_style = "[green]" if state == "ACTIVE" else "[yellow]" | ||
|
|
||
| # Validate timestamp with type and range checks, not just truthiness | ||
| if isinstance(updated, datetime): | ||
| now = datetime.now(tz=UTC) | ||
| # Make tz-aware for comparison (assume UTC if naive) | ||
| ts = updated if updated.tzinfo else updated.replace(tzinfo=UTC) | ||
| if ts > now: | ||
| updated_display = f"[red]{updated} (future timestamp)[/red]" | ||
| else: | ||
| updated_display = str(updated) | ||
| elif updated is not None: | ||
| updated_display = f"[yellow]{updated} (not a datetime)[/yellow]" | ||
| else: | ||
| updated_display = "[yellow]missing[/yellow]" | ||
|
|
||
| table.add_row( | ||
| str(entity_id), | ||
| f"{state_style}{state}[/{state_style.strip('[')}", | ||
| str(ntype or ""), | ||
| updated_display, |
There was a problem hiding this comment.
🧩 Analysis chain
🏁 Script executed:
head -500 src/omnibase_infra/cli/infra_test/verify.py | tail -100Repository: OmniNode-ai/omnibase_infra
Length of output: 3730
🏁 Script executed:
sed -n '430,460p' src/omnibase_infra/cli/infra_test/verify.pyRepository: OmniNode-ai/omnibase_infra
Length of output: 1220
Fix Rich markup closing tag for state styling.
state_style.strip('[') on "[green]" produces "green]", creating an invalid closing tag [/green]]. Use strip('[]') to correctly produce "green", then add the closing ] explicitly to generate valid markup.
🛠️ Proposed fix
- f"{state_style}{state}[/{state_style.strip('[')}",
+ f"{state_style}{state}[/{state_style.strip('[]')}]",🤖 Prompt for AI Agents
In `@src/omnibase_infra/cli/infra_test/verify.py` around lines 433 - 455, The
closing Rich markup tag for state styling is incorrectly built using
state_style.strip('[') which yields "green]" and creates invalid markup; update
the construction in the table.add_row call that builds the state cell (the
variables state_style and state are used there) so you strip both brackets
(e.g., strip('[]')) or otherwise extract the bare style name from state_style
and then concatenate it into a proper closing tag like f"[/{style_name}]" so the
final string is a valid closing tag for the opening state_style.
Summary
onex-infra-testCLI for end-to-end validation of the ONEX registration pipeline using Click + Richintrospectcommand to publish test events via rpk, andverifycommands for Consul, PostgreSQL, and Kafka state validationTest plan
poetry run pytest tests/unit/cli/infra_test/onex-infra-test --helpshows all commandsonex-infra-test env upstarts docker-compose.e2e.yml and waits for healthonex-infra-test verify registrychecks Consul + PostgreSQLonex-infra-test verify topicsvalidates ONEX naming complianceonex-infra-test run --suite smokepasses with live infrastructureSummary by CodeRabbit
New Features
Tests
Chores