Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
93 changes: 74 additions & 19 deletions backend/src/kv-credit-store.ts
Original file line number Diff line number Diff line change
Expand Up @@ -89,17 +89,24 @@ export class KVCreditStore {
if (effectiveBonusUsd > 0n) {
await this.addMonthlyBonusUsed(rootAccount, month, effectiveBonusUsd);
}
// The G$ and USD totals are credited in `markFundingResult`, once funding actually lands. Only
// the outstanding balance moves here, because the entry is at this point owed but not yet paid.
await this.updateUser(account, rootAccount, (current) => ({
...current,
rootAccount: rootAccount,
createdAt: current.createdAt ?? now,
updatedAt: now,
streamFlowRateWeiPerSecond: input.flowRate ? input.flowRate.toString() : current.streamFlowRateWeiPerSecond,
totalGdDepositedWei: addDecimalStrings(current.totalGdDepositedWei, entry.gdAmountWei),
totalGDStreamedWei: input.source.startsWith("stream") ? addDecimalStrings(current.totalGDStreamedWei, entry.gdAmountWei) : current.totalGDStreamedWei,
totalOutstandingFundingUsd: addDecimalStrings(current.totalOutstandingFundingUsd, entry.totalCreditUsd)
}));

// Written separately from the mirrored update above, because the flow rate belongs to this
// account alone -- see `recordStreamFlowRate`. `!== undefined` rather than a truthiness check:
// a terminated stream reports 0, and reading that as "no value given" would leave the closed
// stream's rate on the profile forever.
if (input.flowRate !== undefined) {
await this.recordStreamFlowRate(await this.getUser(account), input.flowRate);
}

logInfo("kv.credit.recorded", {
entryId: entry.id,
account: redactAddress(entry.account),
Expand All @@ -115,6 +122,38 @@ export class KVCreditStore {
return entry;
}

/**
* Sync a profile's recorded flow rate to what the subgraph reports, including 0 for a stream that
* has closed. The subgraph is the source of truth here: the termination event that would
* otherwise write the 0 is pushed in by `POST /v1/celo/events/record` and may never arrive.
*
* Deliberately leaves `lastStreamCreditAt` and every total alone -- this corrects current state,
* it does not settle anything. The caller must pass a freshly read profile: it is both the basis
* for the comparison and the record written back. Updates that object in place to match what was
* stored, and returns whether a write was needed.
*/
async recordStreamFlowRate(profile: UserCreditProfile, flowRate: bigint): Promise<boolean> {
const next = flowRate.toString();
if (profile.streamFlowRateWeiPerSecond === next) return false;
const now = new Date().toISOString();
// Writes this profile only -- deliberately not through `updateUser`, which mirrors every change
// onto the GoodID root. The totals mirror correctly because they accumulate; a flow rate is an
// absolute value, so mirroring it would leave the root holding whichever sub-account wrote last.
// The root's rate is maintained separately, as the sum across the identity's accounts.
const previousFlowRateWeiPerSecond = profile.streamFlowRateWeiPerSecond;
// Keep the caller's copy in step with what was stored, so a profile held across a loop and
// logged later does not report a stale rate or timestamp.
profile.streamFlowRateWeiPerSecond = next;
profile.updatedAt = now;
await this.putJson(`${USER_PREFIX}${profile.account}`, profile);
logInfo("kv.stream.flow-rate-synced", {
account: redactAddress(profile.account),
previousFlowRateWeiPerSecond,
flowRateWeiPerSecond: next
});
return true;
}

async markFundingResult(
entry: GdCreditEntry,
result: { funded: boolean; id?: string; txHash?: string; error?: string; credited?: boolean }
Expand All @@ -133,22 +172,27 @@ export class KVCreditStore {
entry.fundingError = result.error;
await this.putJson(`${GD_CREDIT_PREFIX}${entry.id}`, entry);

if (result.funded) {
const now = new Date().toISOString();
const credited = result.credited !== false;
await this.updateUser(entry.account, entry.rootAccount, (current) => {
const outstanding = BigInt(current.totalOutstandingFundingUsd);
const creditAmount = BigInt(entry.totalCreditUsd);
return {
...current,
updatedAt: now,
lastStreamCreditAt: entry.source.startsWith("stream") ? now : current.lastStreamCreditAt,
totalPrincipalUsd: (BigInt(current.totalPrincipalUsd) + (credited ? BigInt(entry.principalUsd) : 0n)).toString(),
totalBonusUsd: (BigInt(current.totalBonusUsd) + (credited ? BigInt(entry.bonusUsd) : 0n)).toString(),
totalOutstandingFundingUsd: (outstanding > creditAmount ? outstanding - creditAmount : 0n).toString()
};
});
}
// `funded` and `failed` are both terminal -- funding is never retried for an entry -- so either
// way the entry stops being outstanding. Only a funded entry adds to the lifetime totals; a
// failed one must leave them untouched, or the G$ counters drift above the credit actually
// granted (what produced a 6.2M G$ "deposited" profile against a $4 balance).
const now = new Date().toISOString();
const credited = result.funded && result.credited !== false;
await this.updateUser(entry.account, entry.rootAccount, (current) => {
const outstanding = BigInt(current.totalOutstandingFundingUsd);
const creditAmount = BigInt(entry.totalCreditUsd);
return {
...current,
updatedAt: now,
lastStreamCreditAt: result.funded && entry.source.startsWith("stream") ? now : current.lastStreamCreditAt,
totalGdDepositedWei: credited ? addDecimalStrings(current.totalGdDepositedWei, entry.gdAmountWei) : current.totalGdDepositedWei,
totalGDStreamedWei:
credited && entry.source.startsWith("stream") ? addDecimalStrings(current.totalGDStreamedWei, entry.gdAmountWei) : current.totalGDStreamedWei,
Comment on lines +188 to +190

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@sirpy
Valid. This race already exists on main for the USD totals;
The proper fix is to split per-wallet profiles from the identity aggregate.
I do not fix this in the pr

totalPrincipalUsd: (BigInt(current.totalPrincipalUsd) + (credited ? BigInt(entry.principalUsd) : 0n)).toString(),
totalBonusUsd: (BigInt(current.totalBonusUsd) + (credited ? BigInt(entry.bonusUsd) : 0n)).toString(),
totalOutstandingFundingUsd: (outstanding > creditAmount ? outstanding - creditAmount : 0n).toString()
};
});
logInfo("kv.funding.result", {
entryId: entry.id,
account: redactAddress(entry.account),
Expand Down Expand Up @@ -209,6 +253,17 @@ export class KVCreditStore {
};
}

/**
* The stored profile, or undefined when the account has none yet. `getUser` synthesizes a default
* for unknown accounts, which cannot distinguish "never seen" from "seen, all zero" -- callers
* that must not create a profile as a side effect need this instead.
*/
async getSavedUser(account: string): Promise<UserCreditProfile | undefined> {
const normalized = normalizeAccount(account);
const saved = await this.getJson<Partial<UserCreditProfile>>(`${USER_PREFIX}${normalized}`);
return saved ? normalizeProfile(saved, normalized) : undefined;
}

async getUser(account: string): Promise<UserCreditProfile> {
const normalized = normalizeAccount(account);
const saved = await this.getJson<Partial<UserCreditProfile>>(`${USER_PREFIX}${normalized}`);
Expand Down
Loading
Loading