Skip to content

Commit 0e13264

Browse files
authored
Merge pull request #472 from QuantStrategyLab/codex/reconciliation-request-correlation-20260905
fix: correlate IBKR reconciliation receipts
2 parents 8ae6c46 + 0145a4d commit 0e13264

6 files changed

Lines changed: 106 additions & 11 deletions

File tree

.github/workflows/collect-reconciliation-evidence.yml

Lines changed: 14 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -93,6 +93,7 @@ jobs:
9393
set -euo pipefail
9494
test -n "$SERVICE_URL"
9595
job_name="ibkr-reconcile-${PROFILE}-${GITHUB_RUN_ID}-${GITHUB_RUN_ATTEMPT}"
96+
request_id="$(python -c 'import uuid; print(uuid.uuid4())')"
9697
requested_at="$(date -u +%Y-%m-%dT%H:%M:%SZ)"
9798
case "$PROFILE" in
9899
soxl_soxx_trend_income) delay_minutes=3 ;;
@@ -116,6 +117,7 @@ jobs:
116117
--time-zone 'Etc/UTC' \
117118
--uri "${SERVICE_URL}/reconcile" \
118119
--http-method POST \
120+
--headers "X-QSL-Reconciliation-Request-Id=${request_id}" \
119121
--oidc-service-account-email "$GCP_SCHEDULER_SERVICE_ACCOUNT" \
120122
--oidc-token-audience "$SERVICE_URL" \
121123
--attempt-deadline '300s' \
@@ -124,6 +126,7 @@ jobs:
124126
--quiet >/dev/null
125127
{
126128
echo "job_name=$job_name"
129+
echo "request_id=$request_id"
127130
echo "requested_at=$requested_at"
128131
} >> "$GITHUB_OUTPUT"
129132
@@ -135,19 +138,19 @@ jobs:
135138
env:
136139
PROFILE: ${{ matrix.profile }}
137140
SERVICE: ${{ matrix.service }}
138-
SERVICE_URL: ${{ steps.audience.outputs.service_url }}
141+
REQUEST_ID: ${{ steps.scheduler.outputs.request_id }}
139142
REQUESTED_AT: ${{ steps.scheduler.outputs.requested_at }}
140143
run: |
141144
set -euo pipefail
142-
test -n "$SERVICE_URL"
145+
test -n "$REQUEST_ID"
143146
test -n "$REQUESTED_AT"
144147
mkdir -p reports
145148
log_entry=''
146149
logging_read_failures=0
147-
for attempt in $(seq 1 96); do
150+
for _ in $(seq 1 96); do
148151
logging_output=''
149152
logging_read_status=0
150-
logging_output="$(gcloud logging read "resource.type=\"cloud_run_revision\" AND resource.labels.service_name=\"${SERVICE}\" AND timestamp>=\"${REQUESTED_AT}\" AND textPayload:\"execution_report gs://\"" --project "$GCP_PROJECT_ID" --freshness=15m --limit=10 --format=json 2>/dev/null)" || logging_read_status=$?
153+
logging_output="$(gcloud logging read "resource.type=\"cloud_run_revision\" AND resource.labels.service_name=\"${SERVICE}\" AND timestamp>=\"${REQUESTED_AT}\" AND textPayload:\"reconciliation_receipt_ready request_id=${REQUEST_ID} report_uri=gs://\"" --project "$GCP_PROJECT_ID" --freshness=15m --limit=10 --format=json 2>/dev/null)" || logging_read_status=$?
151154
if [ "$logging_read_status" -ne 0 ]; then
152155
logging_read_failures=$((logging_read_failures + 1))
153156
sleep 5
@@ -157,7 +160,7 @@ jobs:
157160
echo "reconciliation_log_query_failed class=invalid_logging_json" >&2
158161
exit 1
159162
fi
160-
log_entry="$(jq -c 'first(.[] | select(.textPayload? | startswith("execution_report gs://")) | select(.resource.labels.revision_name? | type == "string" and length > 0)) // empty' <<< "$logging_output")"
163+
log_entry="$(jq -c --arg request_id "$REQUEST_ID" 'first(.[] | select(.textPayload? | type == "string") | select(.textPayload | startswith("reconciliation_receipt_ready request_id=" + $request_id + " report_uri=gs://")) | select(.resource.labels.revision_name? | type == "string" and length > 0)) // empty' <<< "$logging_output")"
161164
if [ -n "$log_entry" ]; then
162165
break
163166
fi
@@ -173,16 +176,18 @@ jobs:
173176
fi
174177
exit 1
175178
fi
176-
report_uri="$(jq -er '.textPayload | capture("^execution_report (?<uri>gs://[^[:space:]]+)$").uri' <<< "$log_entry")"
179+
receipt="$(jq -cer --arg request_id "$REQUEST_ID" '.textPayload | capture("^reconciliation_receipt_ready request_id=(?<request_id>[0-9a-f-]{36}) report_uri=(?<uri>gs://[^[:space:]]+)$") | select(.request_id == $request_id)' <<< "$log_entry")"
180+
report_uri="$(jq -er '.uri' <<< "$receipt")"
177181
serving_revision="$(jq -er '.resource.labels.revision_name | select(type == "string" and length > 0)' <<< "$log_entry")"
178182
revision_labels="$(gcloud run revisions describe "$serving_revision" --project "$GCP_PROJECT_ID" --region "$GCP_REGION" --format='json(metadata.labels)')"
179183
service_revision_commit_sha="$(jq -er '."commit-sha" | select(type == "string" and test("^[0-9a-f]{40}$"))' <<< "$revision_labels")"
180184
service_deploy_run_id="$(jq -er '."github-run-id" | select(type == "string" and test("^[0-9]+$"))' <<< "$revision_labels")"
181-
gcloud storage cat "$report_uri" | jq --arg profile "$PROFILE" '
185+
gcloud storage cat "$report_uri" | jq --arg profile "$PROFILE" --arg request_id "$REQUEST_ID" '
182186
.diagnostics.broker_reconciliation as $candidate
183187
| {
184188
schema_version: "ibkr_reconciliation_artifact.v1",
185189
strategy_profile: $profile,
190+
reconciliation_request_id: .diagnostics.reconciliation_request_id,
186191
report_status: .status,
187192
reconciliation: $candidate,
188193
errors: [
@@ -194,9 +199,10 @@ jobs:
194199
]
195200
}
196201
' > "reports/${PROFILE}.json"
197-
jq -e --arg profile "$PROFILE" '
202+
jq -e --arg profile "$PROFILE" --arg request_id "$REQUEST_ID" '
198203
.schema_version == "ibkr_reconciliation_artifact.v1"
199204
and .strategy_profile == $profile
205+
and .reconciliation_request_id == $request_id
200206
and .reconciliation.schema_version == "ibkr_reconciliation_candidate.v1"
201207
and .reconciliation.evidence.platform_id == "ibkr"
202208
and .reconciliation.evidence.strategy_profile == $profile

application/reconciliation_reporting.py

Lines changed: 21 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,7 @@
33
from __future__ import annotations
44

55
from collections.abc import Mapping
6+
import re
67
from typing import Any
78

89
_CANDIDATE_FIELDS = (
@@ -52,6 +53,19 @@
5253
"started_at",
5354
"finished_at",
5455
)
56+
_RECONCILIATION_REQUEST_ID_PATTERN = re.compile(
57+
r"[0-9a-f]{8}-(?:[0-9a-f]{4}-){3}[0-9a-f]{12}"
58+
)
59+
60+
61+
def normalize_reconciliation_request_id(value: object) -> str | None:
62+
"""Return a safe opaque correlation ID or reject the untrusted value."""
63+
if not isinstance(value, str):
64+
return None
65+
normalized = value.strip().lower()
66+
if _RECONCILIATION_REQUEST_ID_PATTERN.fullmatch(normalized):
67+
return normalized
68+
return None
5569

5670

5771
def _selected_mapping(value: object, fields: tuple[str, ...]) -> dict[str, Any]:
@@ -89,6 +103,13 @@ def build_persistable_reconciliation_report(report: Mapping[str, object]) -> dic
89103
)
90104
if candidate:
91105
safe_report["diagnostics"]["broker_reconciliation"] = candidate
106+
request_id = normalize_reconciliation_request_id(
107+
(report.get("diagnostics") or {}).get("reconciliation_request_id")
108+
if isinstance(report.get("diagnostics"), Mapping)
109+
else None
110+
)
111+
if request_id:
112+
safe_report["diagnostics"]["reconciliation_request_id"] = request_id
92113
failure = _selected_mapping(
93114
(report.get("diagnostics") or {}).get("broker_reconciliation_failure")
94115
if isinstance(report.get("diagnostics"), Mapping)

main.py

Lines changed: 15 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -28,7 +28,10 @@
2828
build_reconciliation_candidate,
2929
collect_read_only_reconciliation_observations,
3030
)
31-
from application.reconciliation_reporting import build_persistable_reconciliation_report
31+
from application.reconciliation_reporting import (
32+
build_persistable_reconciliation_report,
33+
normalize_reconciliation_request_id,
34+
)
3235
from application.runtime_broker_adapters import (
3336
IBKRGatewayUnavailableError,
3437
IBKRTradingPermissionError,
@@ -1882,9 +1885,14 @@ def _handle_reconciliation():
18821885
ib = None
18831886
log_context = None
18841887
report = None
1888+
reconciliation_request_id = normalize_reconciliation_request_id(
1889+
request.headers.get("X-QSL-Reconciliation-Request-Id")
1890+
)
18851891
try:
18861892
log_context = build_request_log_context()
18871893
report = build_execution_report(log_context, dry_run_only_override=True)
1894+
if reconciliation_request_id is not None:
1895+
report.setdefault("diagnostics", {})["reconciliation_request_id"] = reconciliation_request_id
18881896
runtime_target = RUNTIME_SETTINGS.runtime_target
18891897
if runtime_target is None:
18901898
raise IBKRReconciliationReadError(
@@ -1999,6 +2007,12 @@ def _handle_reconciliation():
19992007
and not any(character.isspace() for character in report_path)
20002008
):
20012009
print(f"execution_report {report_path}", flush=True)
2010+
if reconciliation_request_id is not None:
2011+
print(
2012+
"reconciliation_receipt_ready "
2013+
f"request_id={reconciliation_request_id} report_uri={report_path}",
2014+
flush=True,
2015+
)
20022016
else:
20032017
print("broker reconciliation report persisted", flush=True)
20042018
except Exception as persist_exc:

tests/test_reconciliation_evidence_workflow.py

Lines changed: 18 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -82,3 +82,21 @@ def test_reconciliation_evidence_preserves_sanitized_log_query_failure_class() -
8282
assert 'reconciliation_log_query_timeout class=no_matching_candidate' in workflow
8383
assert 'gcloud logging read' in workflow
8484
assert '--format=json 2>/dev/null' in workflow
85+
86+
87+
def test_reconciliation_evidence_binds_scheduler_request_to_exact_receipt() -> None:
88+
workflow = WORKFLOW.read_text(encoding="utf-8")
89+
90+
required_markers = (
91+
"request_id=\"$(python -c 'import uuid; print(uuid.uuid4())')\"",
92+
"X-QSL-Reconciliation-Request-Id=${request_id}",
93+
"echo \"request_id=$request_id\"",
94+
"REQUEST_ID: ${{ steps.scheduler.outputs.request_id }}",
95+
"reconciliation_receipt_ready request_id=${REQUEST_ID} report_uri=gs://",
96+
"--arg request_id \"$REQUEST_ID\"",
97+
".reconciliation_request_id == $request_id",
98+
)
99+
for marker in required_markers:
100+
assert marker in workflow
101+
102+
assert 'textPayload:\"execution_report gs://\"' not in workflow

tests/test_reconciliation_reporting.py

Lines changed: 24 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -74,6 +74,7 @@ def test_persistable_reconciliation_report_keeps_only_safe_evidence(tmp_path) ->
7474
},
7575
"diagnostics": {
7676
"broker_reconciliation": candidate,
77+
"reconciliation_request_id": "853a2e08-9396-4fe8-89ee-59fb17e40a1d",
7778
"ib_gateway_host": marker,
7879
},
7980
"artifacts": {"strategy_config_path": marker},
@@ -93,6 +94,7 @@ def test_persistable_reconciliation_report_keeps_only_safe_evidence(tmp_path) ->
9394

9495
assert result.local_path is not None
9596
assert payload["diagnostics"]["broker_reconciliation"] == candidate
97+
assert payload["diagnostics"]["reconciliation_request_id"] == "853a2e08-9396-4fe8-89ee-59fb17e40a1d"
9698
assert payload["errors"] == [
9799
{
98100
"stage": "broker_reconciliation",
@@ -111,3 +113,25 @@ def test_persistable_reconciliation_report_keeps_only_safe_evidence(tmp_path) ->
111113
assert "runtime_target" not in payload
112114
assert "account_scope" not in payload
113115
assert "service_name" not in payload
116+
117+
118+
def test_persistable_reconciliation_report_rejects_untrusted_request_id() -> None:
119+
report = {
120+
"schema_version": "runtime_report.v1",
121+
"platform": "ibkr",
122+
"deploy_target": "cloud_run",
123+
"strategy_profile": "global_etf_rotation",
124+
"run_id": "run-001",
125+
"run_source": "cloud_run",
126+
"status": "ok",
127+
"started_at": "2026-09-04T00:00:00Z",
128+
"finished_at": "2026-09-04T00:01:00Z",
129+
"summary": {},
130+
"diagnostics": {"reconciliation_request_id": "unsafe value\nreport_uri=gs://x"},
131+
"artifacts": {},
132+
"errors": [],
133+
}
134+
135+
sanitized = build_persistable_reconciliation_report(report)
136+
137+
assert "reconciliation_request_id" not in sanitized["diagnostics"]

tests/test_request_handling.py

Lines changed: 14 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1446,7 +1446,11 @@ def persist_report(report, **_kwargs):
14461446

14471447
monkeypatch.setattr(strategy_module, "persist_execution_report", persist_report)
14481448

1449-
with strategy_module.app.test_request_context("/reconcile", method="POST"):
1449+
with strategy_module.app.test_request_context(
1450+
"/reconcile",
1451+
method="POST",
1452+
headers={"X-QSL-Reconciliation-Request-Id": "853a2e08-9396-4fe8-89ee-59fb17e40a1d"},
1453+
):
14501454
body, status = strategy_module.handle_reconciliation()
14511455

14521456
assert (body, status) == ("Error", 503)
@@ -1457,6 +1461,9 @@ def persist_report(report, **_kwargs):
14571461
"failure_category": "broker_reconciliation",
14581462
}
14591463
]
1464+
assert observed["report"]["diagnostics"]["reconciliation_request_id"] == (
1465+
"853a2e08-9396-4fe8-89ee-59fb17e40a1d"
1466+
)
14601467
failed_event = observed["events"][-1]
14611468
assert failed_event[0] == "broker_reconciliation_failed"
14621469
assert "error_message" not in failed_event[1]
@@ -1465,7 +1472,12 @@ def persist_report(report, **_kwargs):
14651472
assert "demo-account" not in serialized
14661473
output = capsys.readouterr().out
14671474
assert marker not in output
1468-
assert output == "execution_report gs://private-reports/reconciliation.json\n"
1475+
assert output == (
1476+
"execution_report gs://private-reports/reconciliation.json\n"
1477+
"reconciliation_receipt_ready "
1478+
"request_id=853a2e08-9396-4fe8-89ee-59fb17e40a1d "
1479+
"report_uri=gs://private-reports/reconciliation.json\n"
1480+
)
14691481

14701482

14711483
def test_handle_reconciliation_hides_disconnect_and_persistence_error_details(

0 commit comments

Comments
 (0)