-
Notifications
You must be signed in to change notification settings - Fork 0
Architecture
A decision-reasoning observability system for multi-step pipelines. Captures why decisions were made (candidates, filters, scores, reasoning), not just what functions ran.
- Database Schema
- Type Locations
- SDK Components
- Payload Summarization
- Payload Externalization
- Event Structures
- Complete Data Flow
- Configuration
The system uses 3 tables with clear relationships:
Link: https://dbdiagram.io/d/XRAY-DIAGRAM-6960bf2fd6e030a0248dae71

| Field | Formula | Notes |
|---|---|---|
removed_ratio |
(input_count - output_count) / NULLIF(input_count, 0) |
Computed at query time |
duration_ms |
Stored in Step table | Computed in SDK at step end |
Located in shared/types.py:
class StepType(str, Enum):
filter = "filter" # Removes candidates based on criteria
rank = "rank" # Orders candidates by score
llm = "llm" # LLM/AI model call
retrieval = "retrieval" # Fetches data from external source
transform = "transform" # Transforms data format
other = "other" # Catch-all
class RunStatus(str, Enum):
running = "running" # In progress
success = "success" # Completed successfully
error = "error" # Failed with error
class StepStatus(str, Enum):
running = "running"
success = "success"
error = "error"
class DetailLevel(str, Enum):
summary = "summary" # Counts + samples only
full = "full" # Complete data with truncation| Class | Location | Purpose |
|---|---|---|
Run |
sdk/run.py |
Pipeline execution with Step factory |
Step |
sdk/step.py |
Single decision step lifecycle |
PayloadCollector |
sdk/step.py |
Collects large data during summarization |
Transport |
sdk/transport.py |
Async buffered HTTP sender |
XRayConfig |
sdk/config.py |
SDK configuration |
shared/
└── types.py
├── StepType (enum) ─┐
├── RunStatus (enum) │ Used by both SDK and API
├── StepStatus (enum) │
└── DetailLevel (enum) ─┘
sdk/
├── __init__.py # Public exports
├── config.py
│ └── XRayConfig (class)
│ └── load_config (function)
├── transport.py
│ └── Transport (class)
├── run.py
│ └── Run (class)
└── step.py
├── PayloadCollector (class)
├── Step (class)
├── summarize_payload (function)
├── infer_count (function)
├── is_candidate_list (function)
└── extract_candidate (function)
api/ (to be implemented)
├── models.py # SQLAlchemy ORM models
├── schemas.py # Pydantic request/response
├── store.py # Database operations
└── routes.py # FastAPI endpoints
Location: sdk/step.py:31-63
Collects large data during summarization and assigns reference IDs.
class PayloadCollector:
_payloads: dict[str, Any] # {"p-000": [...], "p-001": "..."}
_counter: int # Sequential counter
def add(self, data: Any) -> str
# Stores data, returns "p-000", "p-001", etc.
def get_payloads(self) -> dict[str, Any] | None
# Returns payloads dict or None if emptyLocation: sdk/step.py:331-521
Represents a single decision step within a Run.
class Step:
# Properties
id: str # UUID
run_id: str # Parent run UUID
name: str # "price_filter"
step_type: StepType # filter/rank/llm/etc.
status: StepStatus # running/success/error
# Methods
def attach_reasoning(self, reasoning: dict | str) -> None
# Adds reasoning info to step
def end(self, output: Any, status: StepStatus = success) -> None
# Ends step with output data
def end_with_error(self, error: BaseException | str) -> None
# Ends step with errorLifecycle:
- Created via
Run.start_step()→ sendsstep_startevent - Execute business logic
- Call
step.end(output)orstep.end_with_error(error)→ sendsstep_endevent
Location: sdk/run.py:15-237
Represents a complete pipeline execution.
class Run:
# Properties
id: str # UUID
pipeline_name: str # "recommendation_pipeline"
status: RunStatus # running/success/error
metadata: dict[str, Any] # request_id, user_id, etc.
# Methods
def start_step(self, name, step_type, input_data, metadata) -> Step
# Creates and starts a new Step
def end(self, output=None, status=success) -> None
# Ends the run
def end_with_error(self, error, output=None) -> None
# Ends the run with error
# Context manager support
def __enter__(self) -> Run
def __exit__(self, exc_type, exc_val, exc_tb) -> boolLifecycle:
- Created directly or via
XRayClient.start_run()→ sendsrun_startevent - Create steps via
run.start_step() - Call
run.end()or use context manager → sendsrun_endevent
Location: sdk/transport.py:16-183
Async buffered transport with fail-open semantics.
class Transport:
# Properties
is_started: bool
queue_size: int
# Methods
async def start(self) -> None
# Creates HTTP client, starts background worker
def send(self, event: dict) -> bool
# Non-blocking queue put (drops if full)
async def shutdown(self, timeout: float = 5.0) -> None
# Graceful shutdown with queue drainFail-Open Behavior:
- Queue full → event dropped (logged)
- Network error → logged, retry with backoff
- SDK never crashes the application
# Truncation limits
MAX_STRING_LENGTH = 1024 # Truncate inline strings at 1KB
MAX_DICT_KEYS = 50 # Max keys to extract from dicts
MAX_PAYLOAD_DEPTH = 5 # Max recursion depth
# Externalization thresholds
LARGE_LIST_THRESHOLD = 100 # Lists ≥100 items → externalize
LARGE_STRING_THRESHOLD = 2048 # Strings ≥2KB → externalize
PREVIEW_SIZE = 5 # Items in preview for large lists
STRING_PREVIEW_SIZE = 100 # Chars in preview for large strings
# ID detection
ID_FIELDS = ("id", "_id", "candidate_id", "item_id", "product_id", "doc_id")Returns count for list-like objects:
- Lists, tuples, sets →
len() - Dicts with keys
items,results,data,records,candidates→ length of that value - Otherwise →
None
Returns True if obj is a list of dicts where every dict has an ID field.
Checks ALL items (no sampling) to ensure correctness.
Extracts {id, score, reason} from a candidate dict:
-
ID: looks in
id,_id,candidate_id,item_id,product_id,doc_id -
Score: looks in
score,rank,relevance,confidence,weight -
Reason: looks in
reason,explanation,rationale,why,filter_reason
Main summarization function. Handles all types recursively:
| Input Type | Condition | Output |
|---|---|---|
None |
- | {"_type": "null", "_value": None} |
bool |
- | {"_type": "bool", "_value": True} |
int/float |
- | {"_type": "int", "_value": 42} |
str |
< 2KB | {"_type": "str", "_length": N, "_value": "..."} |
str |
≥ 2KB | {"_type": "str", "_length": N, "_ref": "p-000", "_preview": "..."} |
bytes |
- | {"_type": "bytes", "_length": N} |
list |
is_candidate_list | {"_type": "candidates", "_count": N, "_candidates": [...]} |
list |
< 100 items | {"_type": "list", "_count": N, "_values": [...]} |
list |
≥ 100 items | {"_type": "list", "_count": N, "_ref": "p-000", "_preview": [...]} |
dict |
- | {"_type": "dict", "_key_count": N, "_keys": [...], "_values": {...}} |
object |
- | {"_type": "ClassName", "_id": "..."} |
| any | depth ≥ 5 | {"_type": "...", "_truncated": true} |
Large data is stored separately to keep summaries compact while preserving full data.
In the SDK, during summarize_payload() execution — before data is sent to the API.
| Data Type | Threshold | Small (Inline) | Large (Externalized) |
|---|---|---|---|
| Lists | 100 items | _values: [all items] |
_ref: "p-000" + _preview
|
| Strings | 2KB | _value: "full string" |
_ref: "p-001" + _preview
|
| Candidates | Never | Always inline (id+score+reason) | N/A |
| Layer | What Happens |
|---|---|
| SDK: summarize_payload() | Large data → collector.add(data) → returns "p-000"
|
| SDK: Summary | Contains {"_ref": "p-000", "_preview": [...]}
|
| SDK: Event | Contains both summary AND _payloads: {"p-000": [...]}
|
| API: /ingest | Extracts _payloads, stores in Payload table |
| DB: Step | Stores summary with _ref pointers only |
| DB: Payload | Stores full data, linked by (step_id, ref_id, event_type)
|
| Query Time | Join Step ↔ Payload to get full data |
Events sent from SDK to API via Transport.
{
"event_type": "run_start",
"id": "run-uuid",
"pipeline_name": "recommendation_pipeline",
"status": "running",
"started_at": "2024-01-15T10:30:00Z",
"input_summary": { ... },
"metadata": { "request_id": "req-123", "user_id": "u-456" },
"request_id": "req-123",
"user_id": "u-456",
"environment": "prod",
"_payloads": { "p-000": [...] }
}{
"event_type": "run_end",
"id": "run-uuid",
"status": "success",
"ended_at": "2024-01-15T10:30:05Z",
"output_summary": { ... },
"error_message": null,
"_payloads": { "p-000": [...] }
}{
"event_type": "step_start",
"id": "step-uuid",
"run_id": "run-uuid",
"step_name": "price_filter",
"step_type": "filter",
"index": 0,
"started_at": "2024-01-15T10:30:01Z",
"input_summary": {
"_type": "candidates",
"_count": 500,
"_candidates": [
{"id": "p-1", "score": 0.9, "reason": null},
{"id": "p-2", "score": 0.85, "reason": null}
]
},
"input_count": 500,
"metadata": { "threshold": 100 },
"_payloads": null
}{
"event_type": "step_end",
"id": "step-uuid",
"run_id": "run-uuid",
"status": "success",
"ended_at": "2024-01-15T10:30:02Z",
"duration_ms": 1200,
"output_summary": {
"_type": "candidates",
"_count": 50,
"_candidates": [
{"id": "p-5", "score": 0.95, "reason": "best match"},
{"id": "p-8", "score": 0.92, "reason": null}
]
},
"output_count": 50,
"reasoning": { "threshold": 100, "explanation": "Filtered by price" },
"error_message": null,
"_payloads": null
}sdk:
base_url: http://localhost:8000 # API endpoint
api_key: your-api-key # Auth token (optional)
buffer_size: 1000 # Max events in queue
flush_interval: 5.0 # Seconds between flushes
batch_size: 100 # Events per HTTP request
http_timeout: 30.0 # HTTP timeout in secondsapi:
database_url: postgresql+asyncpg://localhost:5432/xray
debug: falseThe config file is found by searching from the current directory upward:
/home/user/myproject/src/main.py (running here)
↓ searches upward
/home/user/myproject/xray.config.yaml (found!)
-
Counts as columns:
input_count/output_countindexed for efficient filtering -
removed_ratiocomputed at query time: Allows formula changes without migration -
Payload externalization: Large data stored separately, linked by
_ref - All candidate IDs captured: Never truncate candidate lists (essential for debugging)
- Fail-open SDK: Network errors never crash the application
-
Context manager support:
with Run(...) as run:for automatic cleanup