caishen

PDSP — Export Pipeline & In-Database Compute

RISE Framework Specification

Spec ID: 75 Version: 1.0 Document ID: caishen-rise-pdsp-export-v1.0 Last Updated: 2026-08-01 Depends on: 71 (schema), 72 (incremental update), 73 (events)


Creative Intent

What The Export Pipeline Enables:

Desired Outcomes:

  1. Every refreshed series has a current CSV on disk within seconds
  2. Consumers of those files learn about them without polling
  3. Historical exports are compressed; working exports are not
  4. Any consumer’s required format is producible without changing the store

Export Formats

Three formats were in production, each with a distinct consumer.

1. Averaged (mid-price) CSV — the working format

The default. Collapses bid and ask into a mid price.

Date,Open,High,Low,Close,Volume

Open  = (askO + bidO) / 2
High  = (askH + bidH) / 2
Low   = (askL + bidL) / 2
Close = (askC + bidC) / 2

This is the format the analysis and charting tooling consumes. The row cap keeps it small enough to re-write on every refresh.

2. Full compressed CSV — the archive format

Gzip-compressed, indexed by date, and — unlike the averaged form — retaining both sides.

Date,AskOpen,AskHigh,AskLow,AskClose,BidOpen,BidHigh,BidLow,BidClose,Volume
<outputPath>/<INSTRUMENT>-<PAIR>_<TF><suffix>.csv.gz

Used for archival and for bulk consumption by ML pipelines. This is the format that preserves the most information.

3. Broker-platform CSV — the interchange format

The format the broker’s strategy platform imports.

HDR;<instrument>;<dtFrom>;<dtTo>;<timeframe>;1;<pipSize>
DAT;<date>;<askO>;<askH>;<askL>;<askC>;<bidO>;<bidH>;<bidL>;<bidC>;<volume>;<volume>

Date format is unresolved. The procedure’s own comment block documents the broker’s sample as dd.MM.yyyy HH:mm:ss, but the procedure hands Dt to pandas to_csv, which emits ISO-8601. Either the broker accepts both or this export was never round-tripped. A re-implementation must confirm the required format against the target platform’s importer before relying on this section.

Together with format 2, this is why spec 70 insists bid and ask are stored separately. Any re-implementation that collapses to mid at write time can produce neither.


Export Variants

Beyond format, exports vary along two axes: how the window is bounded and how the file is named.

Variant Window Filename
nb Most recent N rows <pov><suffix>.csv
dt_nb Most recent N rows ending at a date <pov><suffix>.csv
pov Whole series <pov><suffix>.csv
pov_outpath Whole series caller-specified path
gz_nb Most recent N rows <pov><suffix>.csv.gz
ic_pov Explicit date range broker format

Six near-identical procedures, differing in a WHERE clause and a filename. Requirement: collapse into one parameterized export operation:

Export(instrument, timeframe,
       format   : avg | full | broker,
       window   : { lastN } | { from, to } | { all },
       compress : none | gzip,
       outputPath, suffix)
    -> { path, rowCount, bytes }

The In-Database Execution Model

This is the architecturally significant part of the export pipeline, and the part that most needs a deliberate decision.

The exports do not run in the application. They run inside the database engine, as Python, via the engine’s external-script facility:

1.  A stored procedure builds a SQL query string for the requested window
2.  It builds a Python script as a string, interpolating the instrument,
    timeframe, suffix, and output path directly into the source text
3.  It executes the Python with the query's result set bound as an input
    dataframe
4.  The Python reads additional configuration from a JSON file on the
    database host's filesystem
5.  The Python writes the CSV (optionally gzipped) to the host's filesystem
6.  The Python then invokes a shell script via a subprocess call, passing
    the POV, instrument, and timeframe

Step 6 is the notification hook (spec 73). Two are in use: hook-avg-updated.sh after an averaged export, and hook-full-updated.sh after a full export.

Why it was built this way

The reasoning is sound and worth stating plainly:

Why it must not be carried forward as-is

Concern Detail
Injection The Python is assembled by string concatenation with parameters interpolated directly into source text. An instrument name containing a quote or newline injects arbitrary Python running with the database engine’s privileges. The same parameters are interpolated into the SQL string. This is remote code execution reachable from any caller who can name an instrument.
Blast radius The database engine can write arbitrary files and spawn arbitrary processes on its host.
Split configuration The procedures read a JSON config from a hardcoded host path (/work/o/zeutil/conf/...). Behaviour therefore depends on files outside the database, outside the application, and outside version control.
Untestable The Python exists only as a string inside a T-SQL procedure. It cannot be linted, unit-tested, or type-checked.
Silent failure Every caller wraps the export in a catch that discards the error. A broken export is invisible.
Duplicated logic Timeframe normalization (m1 → mi1) is re-implemented inside each Python string, in addition to the four places it already exists in the application (spec 70).
Portability Ties the system to one database product’s optional scripting feature.

Required target architecture

Move export to a subscriber of the price-update event (spec 73):

   Refresh commits
        │
        ▼
   SeriesUpdateCompleted event
        │
        ▼
   Export subscriber
        ├── reads the window through the normal data layer
        ├── formats (avg | full | broker)
        ├── writes to the configured artifact location
        └── emits SeriesExported { instrument, timeframe, path, rows, bytes }

Requirements:

If in-database compute is genuinely wanted for a future workload (large aggregations that would be expensive to move), that is a separate architectural decision to make explicitly — with parameterized, version-controlled, tested code — not something to inherit from the export path.


Normalization

Two procedures implement price normalization, and they are worth preserving as specifications even though the implementations are flawed.

Purpose

Map prices onto a common scale so series with very different absolute levels can be compared, overlaid, or fed to models that assume bounded inputs.

sp_get_pov_norm_range — range mapping

Computes the linear map that carries one timeframe’s observed range onto a target scale, using a reference timeframe to establish the bounds:

referenceTimeframe = M1 (monthly)        -- the widest available view
targetMin, targetMax = 1, 100

rMin = MIN(askL) over (instrument, referenceTimeframe)
rMax = MAX(askH) over (instrument, referenceTimeframe)

for the requested timeframe:
    srcMin = MIN(askL), srcMax = MAX(askH)
    normalize(v) = (v - rMin) / (rMax - rMin) * (targetMax - targetMin) + targetMin

returns { srcMin, normalize(srcMin), srcMax, normalize(srcMax) }

The key idea: normalize against the monthly range, not against the window being normalized. A window normalized against itself always spans the full target range, destroying the information that this week was quiet and last week was volatile. Normalizing against the long-horizon range preserves it.

Defect: the source maximum is computed as MAX(askL) where MAX(askH) is intended (sp_get_pov_norm_range__220830.create.sql — @mi and @mx are both taken from AskL), so the source range is understated. The reference bounds rMin/rMax are correct (Min(AskL) / Max(AskH)).

GetChartNormalized — per-bar normalization

Returns bars with additional normalized columns:

normalize(v, min, max) = ROUND((v - min) / (max - min) * 10 * baseUnitSize,
                               precision + 2)

Scale factor and rounding derive from the instrument’s own properties, which is right — normalized output stays at a meaningful precision for that instrument.

Defect 1 — scope: the min/max are computed over the entire table, across every instrument and every timeframe:

set @minX = (select Min(AskO) from dbo.PDSPPrices)   -- no WHERE clause

The commented-out where Instrument = @i on the very next line shows the intent. As written, an index priced in the thousands and a currency pair priced near 1 share one scale, so every forex pair normalizes to approximately zero. The procedure does not do what its name says.

Defect 2 — precision, and it interacts with defect 1: the bounds are declared decimal with no precision or scale, which defaults to decimal(18,0) — zero decimal places. All four bounds are truncated to integers before the arithmetic. For a currency pair, @minX and @maxX would both truncate to the same integer, making (max - min) zero.

The two defects currently mask each other: because the bounds span the whole table, max - min is large enough to be non-zero, so the division survives. Fixing only the WHERE clause would turn this into a divide-by-zero for every forex pair. Both must be fixed together — declare the bounds with explicit precision sufficient for the instrument, or use a floating type.

Required specification

Normalize(instrument, timeframe, window,
          referenceTimeframe = monthly,
          targetRange = [1, 100])

    bounds = { min: MIN(low),  max: MAX(high) }
             over (instrument, referenceTimeframe)      -- per instrument, always

    for each bar in window:
        emit each price component mapped linearly from bounds into targetRange,
        rounded to the instrument's precision + 2

Requirements: bounds are always scoped to a single instrument; use low for the minimum and high for the maximum; make the reference timeframe and target range parameters; return the bounds alongside the data so a consumer can invert the mapping.


In-Database Indicator Computation

Beyond exports, several artifacts compute indicators in SQL — moving averages over window functions, an awesome-oscillator formulation, and the multi-window average views described in spec 71.

The moving averages use symmetric and asymmetric windows:

AVG(midBar) OVER (ORDER BY dt ROWS BETWEEN 7 PRECEDING AND 7 FOLLOWING)
AVG(midBar) OVER (ORDER BY dt ROWS BETWEEN 5 PRECEDING AND 3 FOLLOWING)
AVG(midBar) OVER (ORDER BY dt ROWS BETWEEN 8 PRECEDING AND 5 FOLLOWING)
AVG(midBar) OVER (ORDER BY dt ROWS BETWEEN 13 PRECEDING AND 8 FOLLOWING)

The (5,3), (8,5), (13,8) triple is the Alligator’s three displaced moving averages expressed in SQL.

Requirement: these were exploratory. The indicator layer (spec 05) is where indicator computation belongs, because it is where the displacement semantics can be made explicit and tested. Preserve the formulas — recorded above and in spec 71 — and discard the SQL implementations.

Critical property to carry across: a window including following rows is forward-looking. It cannot be computed for the newest bars and must never feed a live signal without an explicit lag. In a database view this property is invisible; in the indicator layer it must be a declared attribute of the indicator.


Defects To Fix, Not Port

# Defect Location Fix
1 Code injection. Python and SQL both assembled by string interpolation of caller-supplied parameters. RCE with the database engine’s privileges. every sp_csv_* / sp_ic_csv_* E1, E2
2 Normalization scoped to the whole table — no WHERE Instrument. GetChartNormalized Scope per instrument
3 Normalization bounds truncated to integers (declare @minX as decimal → decimal(18,0)). Masked by defect 2; fixing 2 alone yields divide-by-zero for forex. GetChartNormalized Fix with defect 2, never separately
4 MAX(AskL) where MAX(AskH) intended — source range understated. sp_get_pov_norm_range__220830 Use AskH for the maximum
5 Broker export upper bound is a second lower bound. AND (Dt >= @DtPoint1) AND (Dt >= @DtPoint2) — the range predicate returns everything after the later of the two dates. sp_ic_csv_pov_exporter__220905_noreturn Dt <= @DtPoint2
6 Three disagreeing row-cap defaults (3,000 / 4,000 / JSON file). procedure, application, host config E7 — one setting
7 Configuration split across a host filesystem path (/work/o/zeutil/conf/*.json) outside version control. every export procedure E7
8 Export failures swallowed by every caller. avg_Exporting, full_Exporting, UpdateLastFrom E4
9 Timeframe normalization re-implemented inside each Python string, on top of the four application-side copies. every export procedure One normalization at the boundary (spec 70)
10 Six near-identical procedures differing only in a WHERE clause and a filename. all export procedures One parameterized operation

Scheduling and Concurrency

Exports run as background tasks queued after each series commits. The engine awaits all outstanding tasks before exiting (spec 74).

Observed problems:

Requirements:

S4 matters more than it looks: the notification hook fires on completion, and a consumer that reads the file must not race the writer.


Configuration

Setting Purpose Original
export.path Artifact directory /data/exports/csv (procedure default + JSON file)
export.avgSuffix Averaged filename suffix .avg
export.avgRows Averaged row cap 3,000 (procedure) / 4,000 (application constant) / JSON file — three sources, reconcile to one
export.fullRows Full-export row cap 5,000 (procedure), JSON-overridable
export.brokerPath Broker-format directory /data/exports/iccsv
export.concurrency Max concurrent exports absent

Requirement (repeat of spec 71’s deployment note): exports must not share a mount with the database’s data files. In the original both live under /data, so an export loop can exhaust the database’s disk.


Verification

# Scenario Expected
T1 Export a series twice with no intervening change Byte-identical artifacts
T2 Export while a refresh writes the same series Consumer never sees a partial file (S4)
T3 Instrument name containing a quote or newline Handled as data; no injection (E2)
T4 Export fails (disk full) Refresh still succeeds; failure in the run summary (E3, E4)
T5 40-series refresh with exports enabled Concurrency stays within bound (S1)
T6 Broker-format export Both bid and ask present; precision from instrument properties
T7 Normalize a forex pair and an index Each scaled against its own monthly range
T8 Normalized values Invertible from the returned bounds

Traceability

Concept Original artifact
Averaged CSV export gia-mssql/sp_csv_avg_exporter__220624__noreturn.*.sql (8 revisions)
Full gzip export sp_csv_gz_nb_suffix__220624 in PDSDb.script-221031.sql
Broker-format export sp_ic_csv_pov_exporter__220905_noreturn
POV / date-ranged variants sp_csv_pov_exporter__220622*, sp_csv_dt_nb_suffix__220624*, sp_csv_nb_suffix__220624*
Range normalization gia-mssql/sp_get_pov_norm_range__220830.create.sql
Per-bar normalization GetChartNormalized in PDSDb.script-221031.sql
SQL indicator experiments gia-mssql/x-indicators-in-SQL/
Python-in-SQL experiments gia-mssql/x-PDSP-csv_from_py_MSSQL__220622/
ML Services enablement gia-mssql/inst__Enable_Machine_Learning_Services.sql, SpiderDbDatabase/RnPy_SQL2208/
Application-side triggering PDSEngine2203/Program.cs → avg_Exporting, full_Exporting
Stored-procedure bindings PDSP.Data.SQL/Model/, PDSP.Data.SQL/SP/