caishen

PDSP — Incremental Update

RISE Framework Specification

Spec ID: 72 Version: 1.0 Document ID: caishen-rise-pdsp-incremental-v1.0 Last Updated: 2026-08-01 Depends on: 70 (architecture), 71 (schema)

This is the central specification of the PDSP set. Everything else supports it.


Creative Intent

What Incremental Update Enables:

Desired Outcomes:

  1. Running a refresh twice in a row does nothing the second time
  2. The bar that is still forming is always identifiable and always current
  3. A closed bar is written once and never changes again
  4. A long-idle series catches up without a full re-download
  5. Interruption at any point leaves the store consistent and resumable

Core Concept: Bar Completeness

Every bar in the store is in one of two states.

State Meaning Mutability
Forming (isIncompleted = true) The period this bar covers has not ended. Its high, low, close, and volume can still change. Rewritten on every refresh
Closed (isIncompleted = false) The period has ended. The bar is final. Never rewritten

Invariant C1 — At most one forming bar per series. For any (instrument, timeframe), at most one row has isIncompleted = true.

Invariant C2 — The forming bar is the newest bar. If a forming bar exists for a series, no closed bar in that series has a later dt.

Invariant C3 — Closure is monotonic. A bar transitions forming → closed exactly once, and never back.

These three invariants are what make the algorithm work. Every step below either relies on them or maintains them.


The Anchor

The forming bar is the anchor: the single point from which the next refresh must resume.

FindAnchor(instrument, timeframe):
    candidates = PricePoints
                 WHERE instrument = ? AND timeframe = ? AND isIncompleted = true
                 ORDER BY dt DESCENDING
    return candidates.first() or NULL

Because of invariant C1 there is at most one candidate; the ordering is defensive.

Naming note. The original calls this GetOldest_ToUpdate while ordering descending. It is not a bug — “oldest” means the oldest point in time from which work is required, which is the newest incomplete bar. The name confuses every reader. Call it FindAnchor or FindResumePoint.

Why the forming bar and not simply “the latest bar”: the latest bar may be closed (the market closed between refreshes, or the previous run completed cleanly on a period boundary). Resuming from the latest closed bar would re-request a period already final. Resuming from the anchor requests exactly the periods that can still change, plus the ones that have not been seen at all.


Algorithm A — Anchored Incremental Merge

The primary path. Given a freshly fetched history that overlaps the store, merge it in.

Precondition: freshHistory is ordered ascending by dt and its earliest bar is at or before the anchor’s dt.

UpdateFrom(instrument, timeframe, freshHistory):

    anchor = FindAnchor(instrument, timeframe)
    IF anchor IS NULL:
        RETURN Bootstrap(instrument, timeframe, freshHistory)     -- Algorithm C

    properties = anchor.instrumentProperty
    lastIndex  = freshHistory.count - 1

    resumed = false
    touched = []

    FOR index, bar IN freshHistory:

        -- 1. Skip everything before the anchor. Nothing there can have changed.
        IF bar.dt == anchor.dt:
            resumed = true
        IF NOT resumed:
            CONTINUE

        -- 2. Derive identity from content (spec 70)
        key = MakeBarKey(instrument, timeframe, bar.dt)

        -- 3. Merge
        IF key == anchor.key:
            -- The anchor itself: rewrite in place.
            anchor.overwriteFrom(bar, instrument, timeframe, properties)
            -- It closes iff a later bar exists in the fresh data.
            IF index < lastIndex:
                anchor.isIncompleted = false
            target = anchor
        ELSE:
            -- Everything after the anchor is new.
            target = NewBar(bar, instrument, timeframe, properties)
            target.isIncompleted = (index == lastIndex)
            store.insert(target)

        touched.append(target)

    store.commit()
    RETURN touched

Why each step is the way it is

Step 1 — skip until the anchor. Bars older than the anchor are closed (invariants C2, C3). Re-writing them would be pure cost and would risk overwriting a correct historical bar with a provider’s later revision. The skip is what makes the operation cheap: a fetch of 300 bars where only 2 are new does 2 units of work.

Step 3, closure rule. The anchor closes if and only if the fresh data contains a bar after it. That is the only reliable evidence that its period ended — more reliable than consulting a clock, because it comes from the data provider’s own view of period boundaries. This is the single most important line in the algorithm.

Step 3, the new tail. Exactly one bar in the result is forming: the last one. This re-establishes invariants C1 and C2 for the next run.

Worked example

Store holds ... 12:00 (closed), 16:00 (forming) for EUR/USD H4. Fresh fetch returns 08:00, 12:00, 16:00, 20:00.

Bar Action Resulting state
08:00 skipped (before anchor) unchanged, closed
12:00 skipped (before anchor) unchanged, closed
16:00 anchor — rewritten; a later bar exists → closes updated, closed
20:00 new — inserted; is last → forming inserted, forming

Net: 1 update, 1 insert, 2 skips. Next run’s anchor is 20:00.

Idempotence

Running immediately again: fresh fetch returns the same four bars, anchor is now 20:00. Bars 08:00–16:00 are skipped; 20:00 is rewritten with identical values and stays forming (no later bar). Zero net change. ✅


Algorithm B — Range Fast-Path

When the fetched window does not overlap the store at all, per-bar merging is wasted work: every bar is new, and every existence check is a guaranteed miss.

BulkInsert(instrument, timeframe, freshHistory):

    -- Cheap probe: are ANY period timestamps already present in this window?
    present = store.GetPeriodTimestamps(instrument, timeframe,
                                        freshHistory.first.dt,
                                        freshHistory.last.dt)

    IF present IS EMPTY:
        bars = [ NewBar(b, instrument, timeframe, properties)
                 FOR b IN freshHistory ]
        bars.last.isIncompleted = true

        TRY:
            store.insertRange(bars)                    -- single bulk operation
        CATCH conflict:
            -- Fall back to per-bar: the probe raced, or the window
            -- straddles an existing boundary.
            FOR bar IN bars:
                IF NOT store.exists(bar.key):
                    store.insert(bar)
        RETURN

    ELSE:
        RETURN UpdateFrom(instrument, timeframe, freshHistory)   -- Algorithm A

The probe reads timestamps only (spec 71, vPDSPPricesDt), never bar bodies. On a cold series it turns thousands of individual existence checks into one query plus one bulk insert.

The fallback is mandatory, not defensive. The probe and the insert are not atomic in the original. Preserve the fallback even if the target store offers a native upsert — a bulk insert that partially fails must degrade to per-row rather than abort the run.


Algorithm C — Bootstrap

When FindAnchor returns null the series is either empty or every bar in it is closed.

Bootstrap(instrument, timeframe, freshHistory):
    properties = ResolveInstrumentProperties(instrument)   -- lazy upsert, spec 70
    bars = [ NewBar(b, instrument, timeframe, properties) FOR b IN freshHistory ]
    bars.last.isIncompleted = true
    store.insertRangeIgnoringExisting(bars)
    store.commit()

This path does not exist in the original and its absence is a live defect. UpdateLastFrom dereferences the anchor immediately:

var oldestIncompletedToUpdated = GetOldest_ToUpdate(instrument, timeframe);
var prop = oldestIncompletedToUpdated.InstrumentProperty;   // null reference

GetOldest_ToUpdate explicitly returns null when nothing is incomplete, and carries the author’s own note at that point: “TODO Something to Initialize the POV Instead of Updating it”. Algorithm C is that missing initialization.

The consequence in the original is documented below under The market-closed tension — it is the reason a correct improvement had to be reverted.


The Market-Closed Tension

This is the deepest design problem in the original and it must be resolved deliberately in any re-implementation.

What the code does

The last bar of every fetch is marked forming, unconditionally:

currentIncompletedBar.IsIncompleted = true;

What the code tried to do

At three separate sites — once in UpdateLastFrom and twice in the engine — the same block appears, commented out, behind a region marker naming the reason:

The commented body in each:

// if (DTSApp.MarketIsOpen) { bar.IsIncompleted = true; }
// else { bar.IsIncompleted = false;  // completed because we are in the weekend

The intent is correct: when the market is closed, the last bar’s period has ended, so it should be closed, not left forming.

Why it was reverted

Closing the last bar leaves the series with no anchor. On the next run FindAnchor returns null, and — with Algorithm C missing — the update path crashes. So the author reverted to marking the last bar forming unconditionally.

The cost of the workaround

The most recent bar of every series is permanently re-fetched and re-written on every run, forever, even for series whose market has been closed for days. Across dozens of instruments and eight timeframes this is the dominant cost of a refresh cycle, and it is pure waste.

Resolution

Implement Algorithm C, then implement the completeness decision properly:

DetermineCompleteness(bar, index, lastIndex, timeframe, marketCalendar, now):

    IF index < lastIndex:
        RETURN CLOSED           -- a later bar exists: definitive

    periodEnd = PeriodStart(bar.dt, timeframe) + PeriodDuration(timeframe)

    IF now >= periodEnd AND NOT marketCalendar.isOpenDuring(bar.dt, periodEnd):
        RETURN CLOSED           -- period elapsed and market was shut

    IF now >= periodEnd AND marketCalendar.hasClosedSince(periodEnd):
        RETURN CLOSED           -- period elapsed, session has since ended

    RETURN FORMING

With Algorithm C present, a series with no forming bar is a normal, handled state — not a crash. The resume point then becomes:

ResumeFrom(instrument, timeframe):
    anchor = FindAnchor(...)
    IF anchor: RETURN anchor.dt
    latest = GetLast(...)
    IF latest: RETURN latest.dt + PeriodDuration(timeframe)   -- resume after it
    RETURN NULL                                                -- cold: bootstrap

Requirement: a market calendar is a first-class dependency of the price store, not an afterthought. It must know session open/close per market code (FOREX, COMMODITY, INDICE, TREASURY — the classification the original tried and failed to populate, spec 70) and per holiday. Without it, “is this period over?” is unanswerable and the workaround above is the only safe choice.


Disk Snapshot Cache

Before requesting anything from the provider, the engine checks a disk cache of previously fetched windows.

snapshot = { instrument, timeframe, dtFrom, dtTo, bars[] }
AcquireHistory(instrument, timeframe, dtFrom, dtTo):
    snapshot = diskCache.load(instrument, timeframe, dtFrom, dtTo)
    IF snapshot exists AND snapshot.bars is non-empty:
        RETURN snapshot.bars                       -- no provider call

    bars = provider.fetch(instrument, timeframe, dtFrom, dtTo)

    IF bars is empty:
        wait, retry once
        IF still empty: RETURN empty               -- caller skips this series

    IF NOT resumingFromAnchor:
        diskCache.save(snapshot(bars, dtFrom, dtTo, instrument, timeframe))

    RETURN bars

Purpose. The provider is rate-limited and frequently unavailable. Backfilling years of minute bars means thousands of windowed requests; a crash mid-backfill must not discard the work. The snapshot cache makes a re-run of a historical backfill entirely offline.

The resumingFromAnchor guard. Snapshots are only written for historical window fetches, never for incremental tail fetches. A tail fetch contains a forming bar whose values are provisional; caching it would serve stale data on a later run. Preserve this distinction.

Requirement: the cache is keyed by the full request envelope (instrument, timeframe, dtFrom, dtTo). Two overlapping-but-different windows are different entries. Add an explicit fetchedAt and a retention policy — the original has neither, so a snapshot from a provider’s pre-revision data is served indefinitely.

The original’s on-disk contract

Documented for migration, not as a requirement — a re-implementation may choose any artifact format, but must be able to read the existing one if historical backfill caches are to be reused.

Aspect Value
Root _storage (relative to the engine’s working directory)
Entry directory pdsp-json
Archive directory pdsp-json-archives (archiving disabled — the collection opts out)
Format JSON serialization of the snapshot object
Key A range tlid, not the raw request string: MkPovTimerangedTlid(instrument, timeframe, dtFrom, dtTo) = <fileNamePOV>__<tlidFrom>__<tlidTo> — the same content-derived naming as spec 70’s bar identity, extended with a second timestamp

Note the key derivation inherits the timeframe reduction from spec 70, so a snapshot key’s precision matches its timeframe. A rawId string is also computed at the call site (Program.cs:797-800, /→- and .→_) but is dead — the lookup path uses the range tlid.

Defect — read/write key drift. The read passes the loop-local _dtFromString/_dtToString (Program.cs:813-814) while the write constructs the snapshot from the static dtFromString/dtToString (Program.cs:910). The locals are initialized from the statics before the loop, so the two agree on a simple run — but any path that mutates the statics mid-run writes a snapshot under a key the reader will never look up, silently reducing the cache to a write-only store. One request-envelope value, threaded explicitly, removes the whole class of problem.


Full Refresh Cycle

Putting the pieces together, as the engine runs them per (instrument, timeframe):

RefreshSeries(instrument, timeframe, dtFrom, dtTo):

  1.  properties = ResolveInstrumentProperties(instrument)     -- lazy upsert
  2.  bars       = AcquireHistory(instrument, timeframe, dtFrom, dtTo)
  3.  IF bars empty: record skip; RETURN
  4.  bars       = bars ordered ascending by dt
  5.  present    = store.GetPeriodTimestamps(instrument, timeframe,
                                             bars.first.dt, bars.last.dt)
  6.  IF present is empty:  BulkInsert(...)          -- Algorithm B
      ELSE:                 UpdateFrom(...)           -- Algorithm A
                            (or Bootstrap if no anchor -- Algorithm C)
  7.  store.commit()
  8.  emit SeriesUpdated(instrument, timeframe, lastAsk, lastBid)   -- spec 73
  9.  schedule exports asynchronously                              -- spec 75

Steps 8 and 9 are deliberately after the commit and deliberately do not affect the result. A failed hook or a failed export must never fail a refresh.

Batched commits

Within the per-bar path the original commits every N bars rather than once at the end (timeToSave). This bounds memory and makes a long backfill resumable: an interruption loses at most N bars of work, and the next run resumes from the anchor as normal.

Requirement: keep batched commits, make N configurable, and ensure each batch is a real transaction. Partial visibility of a batch is acceptable (readers see a prefix of a series, which is always valid); a torn bar is not.


Ordering Requirement

Everything above depends on freshHistory being ordered ascending by dt. Provider responses are not reliably ordered.

The original sorts explicitly by year, month, day, hour, minute, second in sequence (OrderPriceData). That decomposition is unnecessary — sort by the instant — but the requirement is real:

Requirement O1. Sort ascending by period-start instant before merging. Requirement O2. Deduplicate by derived key before merging. A provider returning the same period twice must not produce two rows; the deterministic key makes this a simple distinct-by operation, and doing it before the merge avoids relying on a storage-level conflict.


Defects To Fix, Not Port

# Defect Location Fix
1 No bootstrap path. Null anchor dereferenced immediately. PDSPPriceComponent.UpdateLastFrom Algorithm C
2 Unconditional forming tail. Last bar always marked forming; correct logic present but commented out at 3 sites. UpdateLastFrom, PDSEngine2203/Program.cs (×2) DetermineCompleteness + market calendar. Depends on fix 1.
3 Existence check is a tautology. PDSPPrices.Any(p => p.PovTlid == p.PovTlid) compares a row to itself — always true if any row exists, so the guarded Add never runs. The return value is also inverted relative to its name. Latent, not live: the method has no callers; the production per-row fallback at Program.cs:1025 uses the correct Any(i => i.PovTlid == povTlidItem). Do not port the method; do not conclude the bulk-insert fallback is broken. PDSEntities2203.AddNewPriceOnly (dead code) Delete, or Any(p => p.PovTlid == candidate.PovTlid) returning whether it was added
4 Anchor matched by exact timestamp equality. bar.Dt.Equals(anchor.Dt) — brittle under the UTC/local ambiguity of spec 70. If no fresh bar matches exactly, resumed never becomes true and the merge silently does nothing. UpdateLastFrom Match on derived key, not on timestamp. Additionally: if no bar matches the anchor, that is an error to report, not a silent no-op.
5 Duplicate market-code logic, both copies non-functional (prop.Instrument == "i"). PDSPInstrumentPropertyComponent, PDSEntities2203 Single correct implementation (spec 70)
6 Export coupled to update. A CSV export runs inside UpdateLastFrom’s commit block. UpdateLastFrom Move behind the event boundary (specs 73 and 75)
7 Silent skip on repeated provider failure. After two empty fetches the series is skipped with a console message only. PDSEngine2203/Program.cs Structured run result (spec 70, Error handling posture)
8 No concurrency guard. Two concurrent refreshes of one series race on the anchor. whole layer Atomic upsert or per-series lease (spec 70, Concurrency model)
9 Completeness is inverted on every copy. SetPDSPPrice(PDSPPrice) and UpdatePDSPPrice(PDSPPrice) pass p.IsIncompleted == null ? false : true — so a bar that is closed (false, not null) is copied as forming (true). This runs on the anchor-overwrite path (PDSPPriceComponent.cs:99) and on the engine’s current-bar update (Program.cs:1091), and it directly violates invariants C1 and C3: a bar that Algorithm A just closed can be re-opened by the next copy. A plausible root cause of the very symptom the “This Add Issues” comments describe. PDSPPrice.cs:85 and :142 Copy the value: p.IsIncompleted ?? false
10 Bar key derived with different flags by writer and reader — m15 can never match. The writer (PDSPEntitiesExtensions.cs:30, PDSEntities2203.cs:298) calls MakePovTlid(i, tf, dt, ReducePovTlid) — four arguments, so _reduce_m15_m5 defaults to false and an m15 key is yyMMddHHmm. The reader (PDSPPriceComponent.cs:92, Program.cs:1062) calls it with five arguments, passing ReducePovTlid (true) into _reduce_m15_m5, producing yyMMddHHm — one character shorter. For m15 only: the anchor comparison at PDSPPriceComponent.cs:95 is never true, so every bar takes the insert path; and the forming-bar lookup at Program.cs:1080 never finds its target. Both dated engine snapshots pass ReducePovTlid_m15_m5 (false) at the same site — the author found this and the fix was never promoted. writer vs. reader call sites One key-derivation function with one flag set, called identically everywhere (spec 70, Bar Identity)

Verification

An implementation is correct when all of the following hold.

Invariants — after any operation, for every series:

Behavioural tests:

# Scenario Expected
T1 Refresh twice with identical provider data Second run: 0 inserts, 0 value changes
T2 Refresh empty store with 300 bars 300 inserts, exactly 1 forming (the last)
T3 Refresh with data extending 2 bars past the anchor 1 update (anchor closes), 2 inserts, 1 forming
T4 Refresh where fresh data ends at the anchor Anchor updated, stays forming, 0 inserts
T5 Refresh a window entirely disjoint from stored data Bulk path taken; probe returns empty
T6 Refresh a window partially overlapping Merge path taken; no duplicate keys
T7 Provider returns bars out of order Result identical to T3
T8 Provider returns a duplicated period One row for that period
T9 Interrupt mid-batch, re-run Converges to the same state as an uninterrupted run
T10 Series where market has been closed 48h With fix 2: last bar closed, next run makes no writes
T11 Fresh data contains no bar matching the anchor Reported as an error, not silently ignored (fix 4)
T12 Two concurrent refreshes of one series No duplicate keys, no lost update (fix 8)

Traceability

Concept Original artifact
Algorithm A (anchored merge) PDSP.Business/PDSPPriceComponent.cs → UpdateLastFrom
Anchor lookup PDSP.Data/Ctx/PDSDataProvider2203Model.PDSEntities2203.cs → GetOldest_ToUpdate
Algorithm B (range fast-path) PDSEngine2203/Program.cs, sureIsNotBecauseRangeCheck branch
Range probe PDSP.Services/PDSPHistoryServices.cs → IsTimerangeInDb; PDSPPriceComponent.GetTimeRange_dev22053114
Per-bar merge + batched commit PDSEngine2203/Program.cs, timeToSave loop
Disk snapshot cache PDSP.Entities/PDSPRawPriceCollection.cs; PDSP.Data/PDSPDJAC.cs
Completeness (commented-out correct form) UpdateLastFrom and PDSEngine2203/Program.cs, #region This Add Issues... blocks
Ordering PDSEntities2203.OrderPriceData