wondering if the current process of obtaining the data makes sense and if more should be done in the postprocessing in terms of truly making sure we get unique values. Whatw ould be the best way of simplifying it further while keeping it filtered as much as possible
wondering if the current process of obtaining the data makes sense and if more should be done in the postprocessing in terms of truly making sure we get unique values. Whatw ould be the best way of simplifying it further while keeping it filtered as much as possible
"""
DPE Migration Smoke Tests and Processor Methods
================================================
This file contains:
1. Smoke test functions (add to impl_amps.py temporarily)
2. Pre-processor methods (add to preprocessor.py)
3. Post-processor methods (add to postprocessor.py)
"""
# =============================================================================
# PART 1: SMOKE TEST FUNCTIONS (impl_amps.py - temporary, delete after testing)
# =============================================================================
def test_recall_quantities():
"""Smoke test for Recalls (RecallQuantities) from GPS3."""
import datetime
import qztable
from qz.data.where import Where, inlist
from equity.dpe import fields
from equity.dpe.shards.trading.stockloan.loans_to_house import Recalls
from equity.dpe.api.impl_amps import DPEAmpsAdapter
from equity.dpe.api.processors.stockloan.preprocessor import recall_quantities_pre_processor
from equity.dpe.api.processors.stockloan.postprocessor import recall_quantities_post_processor
from equity.strats.core.frameworks.wks.predicate import convert
topic_schema = qztable.Schema(
columnNames=['Entity', 'SecurityId', 'PmeId', 'MarketCode', 'Transactions'],
columnTypes=['str', 'str', 'str', 'str', 'pyobj']
)
w = Where(fields.COBDate) == datetime.date.today()
w &= Where(fields.Event) == 2
w &= Where(fields.Context) == "bus"
p = convert(w)
api = DPEAmpsAdapter(
ShardKlass=Recalls,
Group='amps_gmffts_uat',
PublishFormat='json-tcp',
Topic='/GMFFTS/Inventory/GPS/Positions/GpsJson3',
TopicSchema=topic_schema,
PreProcessor=recall_quantities_pre_processor,
PostProcessor=recall_quantities_post_processor,
AdapterOptions=None,
)
result = api.get(Recalls, p)
print(f"Recalls returned {result.qztable.nRows()} rows")
print(result.qztable[:10])
return result.qztable
def test_use_and_available_quantities():
"""Smoke test for Quantities (UseAndAvailableQuantities) from GPS3."""
import datetime
import qztable
from qz.data.where import Where, inlist
from equity.dpe import fields
from equity.dpe.shards.trading.stockloan.loans_to_house import Quantities
from equity.dpe.api.impl_amps import DPEAmpsAdapter
from equity.dpe.api.processors.stockloan.preprocessor import use_and_available_quantities_pre_processor
from equity.dpe.api.processors.stockloan.postprocessor import use_and_available_quantities_post_processor
from equity.strats.core.frameworks.wks.predicate import convert
topic_schema = qztable.Schema(
columnNames=['Entity', 'SecurityId', 'PmeId', 'MarketCode', 'USClassifications'],
columnTypes=['str', 'str', 'str', 'str', 'pyobj']
)
w = Where(fields.COBDate) == datetime.date.today()
w &= Where(fields.Event) == 2
w &= Where(fields.Context) == "bus"
p = convert(w)
api = DPEAmpsAdapter(
ShardKlass=Quantities,
Group='amps_gmffts_uat',
PublishFormat='json-tcp',
Topic='/GMFFTS/Inventory/GPS/Positions/GpsJson3',
TopicSchema=topic_schema,
PreProcessor=use_and_available_quantities_pre_processor,
PostProcessor=use_and_available_quantities_post_processor,
AdapterOptions=None,
)
result = api.get(Quantities, p)
print(f"Quantities returned {result.qztable.nRows()} rows")
print(result.qztable[:10])
return result.qztable
# =============================================================================
# PART 2: PRE-PROCESSOR METHODS (add to preprocessor.py)
# =============================================================================
# File: equity/dpe/api/processors/stockloan/preprocessor.py
"""
Add these imports at the top of preprocessor.py if not already present:
from typing import Type, TypeVar
from equity.dpe import fields
from equity.dpe.api.utils import as_amps
from equity.strats.core.frameworks.wks import Predicate, Shard
S = TypeVar('S', bound=Shard)
"""
def recall_quantities_pre_processor(shardklass, predicate):
"""
Pre-processor for Recalls shard.
Builds AMPS filter for GPS3 topic to retrieve Borrow Recalls.
Releases DPE-specific columns that don't exist in AMPS source.
"""
from equity.dpe import fields
from equity.dpe.api.utils import as_amps
# Release DPE-specific columns that don't exist in AMPS
dpe_specific_cols = [fields.Context, fields.COBDate, fields.Event]
for col in dpe_specific_cols:
if col in predicate.constrained:
predicate = predicate.release(col)
# Build AMPS filter from remaining predicate
amps_filter = as_amps(predicate) if predicate.where else ""
# Base filter for equities from MLPFSDOM/GWIMDOM
base_filter = "/Entity in ('MLPFSDOM', 'GWIMDOM') and /SecurityStatic/SecurityDomain = 'EQ'"
if amps_filter:
return f"({base_filter}) and ({amps_filter})"
return base_filter
def use_and_available_quantities_pre_processor(shardklass, predicate):
"""
Pre-processor for Quantities shard.
Builds AMPS filter for GPS3 topic to retrieve Use and Available quantities.
Releases DPE-specific columns that don't exist in AMPS source.
"""
from equity.dpe import fields
from equity.dpe.api.utils import as_amps
# Release DPE-specific columns that don't exist in AMPS
dpe_specific_cols = [fields.Context, fields.COBDate, fields.Event]
for col in dpe_specific_cols:
if col in predicate.constrained:
predicate = predicate.release(col)
# Build AMPS filter from remaining predicate
amps_filter = as_amps(predicate) if predicate.where else ""
# Base filter for equities from MLPFSDOM/GWIMDOM
base_filter = "/Entity in ('MLPFSDOM', 'GWIMDOM') and /SecurityStatic/SecurityDomain = 'EQ'"
if amps_filter:
return f"({base_filter}) and ({amps_filter})"
return base_filter
# =============================================================================
# PART 3: POST-PROCESSOR METHODS (add to postprocessor.py)
# =============================================================================
# File: equity/dpe/api/processors/stockloan/postprocessor.py
"""
Add these imports at the top of postprocessor.py if not already present:
from typing import Type, TypeVar
from equity.strats.core.frameworks.wks import Predicate, Shard
from equity.strats.core.frameworks.wks.utils import only
from equity.dpe import fields
import qztable
S = TypeVar('S', bound=Shard)
"""
def recall_quantities_post_processor(table, shardklass, predicate):
"""
Post-processor for Recalls shard.
Extracts recall quantities from Transactions array where Classification == 'Borrow Recalls'.
Adds Context and Event fields for DPE compliance.
"""
from equity.strats.core.frameworks.wks.utils import only
from equity.dpe import fields
import qztable
# Get context from predicate
context_type = only(predicate, fields.Context)
event_type = only(predicate, fields.Event) if fields.Event in predicate.constrained else 2
# Remove existing context if present
if fields.Context in table.columnNames():
table = table.project(fields.Context, exclude=True)
# Extract recall quantities from Transactions array
rows = []
for i in range(table.nRows()):
row_data = table[i]
entity = row_data.get('Entity', '')
security_id = row_data.get('SecurityId', '')
pme_id = row_data.get('PmeId', '')
market_code = row_data.get('MarketCode', '')
# Sum quantities from Transactions where Classification == 'Borrow Recalls'
transactions = row_data.get('Transactions', [])
recall_qty = sum(
t.get('Quantity', 0)
for t in transactions
if t.get('Classification') == 'Borrow Recalls'
)
# Only include rows with non-zero recall quantities
if recall_qty != 0:
rows.append([entity, security_id, pme_id, market_code, recall_qty])
# Create result table with proper schema
schema = qztable.Schema(
[fields.Entity, fields.Security, fields.PME, 'MarketCode', fields.Quantity],
['str', 'str', 'str', 'str', 'decimal']
)
result = qztable.Table(schema)
if rows:
result.append(rows)
# Add DPE compliance fields - MANDATORY for cash flow system
result = result.extendConst(context_type, fields.Context, 'string')
result = result.extendConst(event_type, fields.Event, 'int32')
return result
def use_and_available_quantities_post_processor(table, shardklass, predicate):
"""
Post-processor for Quantities shard.
Extracts multiple quantity fields from USClassifications JSON object:
- BAMLDTCStockLoanContractualQty / ActualQty (DTC 5143)
- BAMLDTCHouseContractualQty / ActualQty (DTC 161 MLPFS)
- GWIMDTCHouseContractualQty / ActualQty (DTC 161 GWIM)
Adds Context and Event fields for DPE compliance.
"""
from equity.strats.core.frameworks.wks.utils import only
from equity.dpe import fields
import qztable
# Get context from predicate
context_type = only(predicate, fields.Context)
event_type = only(predicate, fields.Event) if fields.Event in predicate.constrained else 2
# Remove existing context if present
if fields.Context in table.columnNames():
table = table.project(fields.Context, exclude=True)
# Extract quantities from USClassifications JSON object
rows = []
for i in range(table.nRows()):
row_data = table[i]
entity = row_data.get('Entity', '')
security_id = row_data.get('SecurityId', '')
pme_id = row_data.get('PmeId', '')
market_code = row_data.get('MarketCode', '')
us_class = row_data.get('USClassifications', {})
if isinstance(us_class, dict):
rows.append([
entity,
security_id,
pme_id,
market_code,
us_class.get('BAMLDTCStockLoanContractualQty', 0),
us_class.get('BAMLDTCStockLoanActualQty', 0),
us_class.get('BAMLDTCHouseContractualQty', 0),
us_class.get('BAMLDTCHouseActualQty', 0),
us_class.get('GWIMDTCHouseContractualQty', 0),
us_class.get('GWIMDTCHouseActualQty', 0),
])
# Create result table with proper schema
schema = qztable.Schema(
[fields.Entity, fields.Security, fields.PME, 'MarketCode',
'BAMLDTCStockLoanContractualQty', 'BAMLDTCStockLoanActualQty',
'BAMLDTCHouseContractualQty', 'BAMLDTCHouseActualQty',
'GWIMDTCHouseContractualQty', 'GWIMDTCHouseActualQty'],
['str', 'str', 'str', 'str',
'decimal', 'decimal', 'decimal', 'decimal', 'decimal', 'decimal']
)
result = qztable.Table(schema)
if rows:
result.append(rows)
# Add DPE compliance fields - MANDATORY for cash flow system
result = result.extendConst(context_type, fields.Context, 'string')
result = result.extendConst(event_type, fields.Event, 'int32')
return result
# =============================================================================
# MAIN (for testing)
# =============================================================================
def main():
print("Testing Recalls (RecallQuantities)...")
test_recall_quantities()
print("\nTesting Quantities (UseAndAvailableQuantities)...")
test_use_and_available_quantities()
if __name__ == "__main__":
main()