1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
//! Per-tenant quota accounting on `PostgreSQL`.
//!
//! This is the backend the guarantee is actually about. On a single node a
//! ceiling can be held up by almost anything; the moment two instances admit
//! concurrently, only the database can arbitrate — which is why the reservation
//! below is **one statement**, not a count followed by an insert.
use async_trait::async_trait;
use crate::core::{RunId, Spend, StoreError, Timestamp};
use crate::quota::{Halt, HaltScope, QuotaError, QuotaSettlement, QuotaStore};
use super::postgres::{PostgresStore, amount_of, be, pool_err, sql_amount};
#[async_trait]
impl QuotaStore for PostgresStore {
fn tenant(&self) -> &str {
crate::journal::JournalStore::tenant(self)
}
async fn reserve(
&self,
run: RunId,
limit: Option<u32>,
at: Timestamp,
) -> Result<(), QuotaError> {
let mut client = self
.pool_ref()
.get()
.await
.map_err(|e| QuotaError::Unavailable(pool_err(&e).to_string()))?;
let tenant = self.tenant_name();
let run = run.to_string();
let Some(limit) = limit else {
// No ceiling: still record the run, so `running()` answers honestly
// and a ceiling added later starts from the truth. No lock either —
// there is no decision here for two admissions to disagree about.
client
.execute(
"INSERT INTO quota_running (tenant, run_id, admitted_at)
VALUES ($1, $2, $3)
ON CONFLICT (tenant, run_id) DO NOTHING",
&[&tenant, &run, &at.unix_timestamp()],
)
.await
.map_err(|e| QuotaError::Unavailable(be(&e).to_string()))?;
return Ok(());
};
// The count and the insert decide together **under a per-tenant
// advisory lock**, because nothing weaker serialises them. This
// statement's earlier form ran without the lock, on a comment claiming
// the decision happened "inside the row lock the write takes" — and no
// such lock exists: two INSERTs of *different* rows lock nothing in
// common, each count subquery reads its own statement snapshot under
// READ COMMITTED, and two admissions racing for one remaining slot
// both passed `count < limit` and both landed. A ceiling that admits
// limit+k exactly under concurrent load is the catalogued
// yields-under-load shape, on the one control whose whole promise is
// surviving a second instance.
//
// The lock is transaction-scoped and per tenant — admissions for one
// tenant serialise, which is the semantic the ceiling requires; other
// tenants' admissions do not wait. The length prefix keeps
// `("acme", …)` from colliding with a tenant literally named
// `"acme…"` under concatenation.
//
// `ON CONFLICT DO NOTHING` keeps a retried admission idempotent: the
// run already holds its slot, so re-reserving must neither take a
// second nor be refused against a ceiling it is already counted in.
let tx = client
.transaction()
.await
.map_err(|e| QuotaError::Unavailable(be(&e).to_string()))?;
tx.query_one(
"SELECT pg_advisory_xact_lock(hashtextextended($1, 0))",
&[&format!("quota-admission:{}:{tenant}", tenant.len())],
)
.await
.map_err(|e| QuotaError::Unavailable(be(&e).to_string()))?;
let inserted = tx
.execute(
"INSERT INTO quota_running (tenant, run_id, admitted_at)
SELECT $1, $2, $3
WHERE (SELECT COUNT(*) FROM quota_running WHERE tenant = $1) < $4
OR EXISTS (SELECT 1 FROM quota_running
WHERE tenant = $1 AND run_id = $2)
ON CONFLICT (tenant, run_id) DO NOTHING",
&[&tenant, &run, &at.unix_timestamp(), &i64::from(limit)],
)
.await
.map_err(|e| QuotaError::Unavailable(be(&e).to_string()))?;
if inserted == 1 {
tx.commit()
.await
.map_err(|e| QuotaError::Unavailable(be(&e).to_string()))?;
return Ok(());
}
// Nothing was written. Either the tenant is at its ceiling, or the run
// already held a slot and `DO NOTHING` fired — and those are opposite
// answers, so it is read back rather than assumed. Still inside the
// transaction, so the counts are the ones the decision was made from.
let held: i64 = tx
.query_one(
"SELECT COUNT(*) FROM quota_running WHERE tenant = $1 AND run_id = $2",
&[&tenant, &run],
)
.await
.map_err(|e| QuotaError::Unavailable(be(&e).to_string()))?
.get(0);
if held > 0 {
tx.commit()
.await
.map_err(|e| QuotaError::Unavailable(be(&e).to_string()))?;
return Ok(());
}
let running: i64 = tx
.query_one(
"SELECT COUNT(*) FROM quota_running WHERE tenant = $1",
&[&tenant],
)
.await
.map_err(|e| QuotaError::Unavailable(be(&e).to_string()))?
.get(0);
tx.commit()
.await
.map_err(|e| QuotaError::Unavailable(be(&e).to_string()))?;
Err(QuotaError::TooManyRuns {
tenant,
running: u32::try_from(running).unwrap_or(u32::MAX),
})
}
async fn set_halt(&self, scope: &HaltScope, reason: Option<&str>) -> Result<(), StoreError> {
let client = self.pool_ref().get().await.map_err(|e| pool_err(&e))?;
let scope = scope.key();
match reason {
Some(reason) => {
client
.execute(
"INSERT INTO quota_halted (tenant, scope, reason) VALUES ($1, $2, $3)
ON CONFLICT (tenant, scope) DO UPDATE SET reason = EXCLUDED.reason",
&[&self.tenant_name(), &scope, &reason],
)
.await
.map_err(|e| be(&e))?;
}
None => {
client
.execute(
"DELETE FROM quota_halted WHERE tenant = $1 AND scope = $2",
&[&self.tenant_name(), &scope],
)
.await
.map_err(|e| be(&e))?;
}
}
Ok(())
}
async fn halts(&self) -> Result<Vec<Halt>, StoreError> {
let client = self.pool_ref().get().await.map_err(|e| pool_err(&e))?;
let rows = client
.query(
"SELECT scope, reason FROM quota_halted WHERE tenant = $1 ORDER BY scope",
&[&self.tenant_name()],
)
.await
.map_err(|e| be(&e))?;
rows.into_iter()
.map(|row| {
let stored: String = row.get(0);
// Corruption, not a row to skip: a halt this build cannot read
// is one it would run straight through, and from the outside
// that is indistinguishable from a halt that was lifted.
let scope = HaltScope::parse(&stored).ok_or_else(|| StoreError::Corrupt {
seq: 0,
detail: format!(
"quota_halted holds the scope '{stored}', which this build cannot \
read — refusing rather than admitting work an operator stopped"
),
})?;
Ok(Halt {
scope,
reason: row.get(1),
})
})
.collect()
}
async fn release(&self, run: RunId) -> Result<(), StoreError> {
let client = self.pool_ref().get().await.map_err(|e| pool_err(&e))?;
client
.execute(
"DELETE FROM quota_running WHERE tenant = $1 AND run_id = $2",
&[&self.tenant_name(), &run.to_string()],
)
.await
.map_err(|e| be(&e))?;
Ok(())
}
async fn settle(&self, settlement: &QuotaSettlement) -> Result<(), StoreError> {
let mut client = self.pool_ref().get().await.map_err(|e| pool_err(&e))?;
let tx = client.transaction().await.map_err(|e| be(&e))?;
let tenant = self.tenant_name();
let run = settlement.run.to_string();
let epoch = settlement.epoch.cast_signed();
let tokens = sql_amount(settlement.spend.tokens);
let minor_units = sql_amount(settlement.spend.minor_units);
let inserted = tx
.execute(
"INSERT INTO quota_settled
(tenant, run_id, epoch, period, tokens, minor_units, release_slot)
VALUES ($1, $2, $3, $4, $5, $6, $7)
ON CONFLICT (tenant, run_id, epoch) DO NOTHING",
&[
&tenant,
&run,
&epoch,
&settlement.period,
&tokens,
&minor_units,
&settlement.release_slot,
],
)
.await
.map_err(|e| be(&e))?;
if inserted == 0 {
let stored = tx
.query_one(
"SELECT period, tokens, minor_units, release_slot
FROM quota_settled
WHERE tenant = $1 AND run_id = $2 AND epoch = $3",
&[&tenant, &run, &epoch],
)
.await
.map_err(|e| be(&e))?;
let same = stored.get::<_, Option<String>>(0) == settlement.period
&& stored.get::<_, i64>(1) == tokens
&& stored.get::<_, i64>(2) == minor_units
&& stored.get::<_, bool>(3) == settlement.release_slot;
if !same {
return Err(StoreError::Corrupt {
seq: 0,
detail: format!(
"quota pass {run}/{} was settled twice with different payloads",
settlement.epoch
),
});
}
tx.commit().await.map_err(|e| be(&e))?;
return Ok(());
}
if let Some(period) = settlement.period.as_deref() {
// Cast before addition: BIGINT addition can overflow before LEAST
// sees it. The numeric intermediate is exact and the stored total
// saturates at the largest representable non-negative amount.
tx.execute(
"INSERT INTO quota_spent (tenant, period, tokens, minor_units)
VALUES ($1, $2, $3, $4)
ON CONFLICT (tenant, period) DO UPDATE SET
tokens = LEAST(9223372036854775807::numeric,
quota_spent.tokens::numeric + EXCLUDED.tokens::numeric)::bigint,
minor_units = LEAST(9223372036854775807::numeric,
quota_spent.minor_units::numeric + EXCLUDED.minor_units::numeric)::bigint",
&[&tenant, &period, &tokens, &minor_units],
)
.await
.map_err(|e| be(&e))?;
}
if settlement.release_slot {
tx.execute(
"DELETE FROM quota_running WHERE tenant = $1 AND run_id = $2",
&[&tenant, &run],
)
.await
.map_err(|e| be(&e))?;
}
tx.commit().await.map_err(|e| be(&e))?;
Ok(())
}
async fn spent(&self, period: &str) -> Result<Spend, StoreError> {
let client = self.pool_ref().get().await.map_err(|e| pool_err(&e))?;
let row = client
.query_opt(
"SELECT tokens, minor_units FROM quota_spent
WHERE tenant = $1 AND period = $2",
&[&self.tenant_name(), &period.to_owned()],
)
.await
.map_err(|e| be(&e))?;
Ok(row.map_or_else(Spend::default, |r| Spend {
tokens: amount_of(r.get::<_, i64>(0)),
minor_units: amount_of(r.get::<_, i64>(1)),
}))
}
async fn running(&self) -> Result<u32, StoreError> {
let client = self.pool_ref().get().await.map_err(|e| pool_err(&e))?;
let n: i64 = client
.query_one(
"SELECT COUNT(*) FROM quota_running WHERE tenant = $1",
&[&self.tenant_name()],
)
.await
.map_err(|e| be(&e))?
.get(0);
Ok(u32::try_from(n).unwrap_or(u32::MAX))
}
async fn running_runs(&self, limit: usize) -> Result<Vec<RunId>, StoreError> {
let client = self.pool_ref().get().await.map_err(|e| pool_err(&e))?;
let rows = client
.query(
"SELECT run_id FROM quota_running WHERE tenant = $1 \
ORDER BY run_id ASC LIMIT $2",
&[
&self.tenant_name(),
&i64::try_from(limit).unwrap_or(i64::MAX),
],
)
.await
.map_err(|e| be(&e))?;
let mut out = Vec::with_capacity(rows.len());
for row in rows {
let raw: String = row.get(0);
// Damage rather than absence — see the embedded backend for why a
// skipped row is the worst available answer here.
out.push(RunId::parse(&raw).map_err(|e| StoreError::Corrupt {
seq: 0,
detail: format!("bad run id '{raw}' in the quota slot table: {e}"),
})?);
}
Ok(out)
}
}