Skip to content

Commit 98a4860

Browse files
committed
🐛 fix(storage): keep current decisions out of eviction
A policy result was stored once, in the audit log, with the subject index holding only its identifier. Audit rows carry a global serial, so every repository's evaluations interleave in one bounded log, and pruning the oldest row deleted the subject index whenever that row was still the current one. A busy repository could therefore erase a live decision belonging to a different repository, which returned no result for a subject nothing had re-evaluated. Current decisions now hold their own copy of the record, keyed by evaluation serial in a table the audit bound does not apply to. Pruning drops an audit row and nothing else, so no amount of history churn can reach live state. Keying by serial rather than by subject keeps the newest-first scan the artifact lookup relies on, and lets that lookup read one table instead of joining back to history.
1 parent 9447fb6 commit 98a4860

9 files changed

Lines changed: 279 additions & 202 deletions

File tree

crates/peryx-ecosystem-pypi/src/migration.rs

Lines changed: 10 additions & 23 deletions
Original file line numberDiff line numberDiff line change
@@ -56,11 +56,10 @@ impl MetadataMigration for PypiPlugin {
5656
MetadataRecordSet::QuotaReservation => {
5757
rewrite_json::<QuotaReservationRecord, LegacyQuotaReservation, _>(record, Into::into)
5858
}
59-
MetadataRecordSet::PolicyDecisionHistory => {
59+
MetadataRecordSet::PolicyDecisionHistory | MetadataRecordSet::PolicyDecisionCurrentById => {
6060
rewrite_json::<PolicyDecisionRecord, LegacyPolicyDecisionRecord, _>(record, Into::into)
6161
}
62-
MetadataRecordSet::PolicyDecisionCurrent => rewrite_policy_subject(record, true),
63-
MetadataRecordSet::PolicyDecisionCurrentById => rewrite_policy_subject(record, false),
62+
MetadataRecordSet::PolicyDecisionCurrent => rewrite_policy_subject(record),
6463
MetadataRecordSet::Analytics => rewrite_analytics(record),
6564
})
6665
}
@@ -84,30 +83,18 @@ where
8483
})
8584
}
8685

87-
fn rewrite_policy_subject(record: &MetadataRecord, key: bool) -> Option<MetadataRecord> {
88-
let subject = if key {
89-
record.key.as_bytes()
90-
} else {
91-
record.value.as_slice()
92-
};
93-
if serde_json::from_slice::<PolicyDecisionSubject>(subject).is_ok() {
86+
/// Current decisions are keyed by their subject, so this record set migrates in its key.
87+
fn rewrite_policy_subject(record: &MetadataRecord) -> Option<MetadataRecord> {
88+
if serde_json::from_str::<PolicyDecisionSubject>(&record.key).is_ok() {
9489
return None;
9590
}
96-
let Ok(legacy) = serde_json::from_slice::<LegacyPolicyDecisionSubject>(subject) else {
91+
let Ok(legacy) = serde_json::from_str::<LegacyPolicyDecisionSubject>(&record.key) else {
9792
return None;
9893
};
99-
let encoded =
100-
serde_json::to_string(&PolicyDecisionSubject::from(legacy)).expect("owned policy decision subjects serialize");
101-
Some(if key {
102-
MetadataRecord {
103-
key: encoded,
104-
value: record.value.clone(),
105-
}
106-
} else {
107-
MetadataRecord {
108-
key: record.key.clone(),
109-
value: encoded.into_bytes(),
110-
}
94+
Some(MetadataRecord {
95+
key: serde_json::to_string(&PolicyDecisionSubject::from(legacy))
96+
.expect("owned policy decision subjects serialize"),
97+
value: record.value.clone(),
11198
})
11299
}
113100

crates/peryx-ecosystem-pypi/tests/unit/migration_tests.rs

Lines changed: 91 additions & 71 deletions
Original file line numberDiff line numberDiff line change
@@ -152,10 +152,12 @@ fn test_migration_rewrites_legacy_quota_reservations() {
152152
);
153153
}
154154

155-
#[test]
156-
fn test_migration_rewrites_legacy_policy_history() {
155+
#[rstest::rstest]
156+
#[case::audit_log(MetadataRecordSet::PolicyDecisionHistory)]
157+
#[case::current(MetadataRecordSet::PolicyDecisionCurrentById)]
158+
fn test_migration_rewrites_legacy_policy_records(#[case] record_set: MetadataRecordSet) {
157159
let migrated = rewritten(
158-
MetadataRecordSet::PolicyDecisionHistory,
160+
record_set,
159161
&record(
160162
"decision",
161163
json!({
@@ -196,10 +198,8 @@ fn test_migration_rewrites_legacy_policy_history() {
196198
);
197199
}
198200

199-
#[rstest::rstest]
200-
#[case::subject_key(MetadataRecordSet::PolicyDecisionCurrent, true)]
201-
#[case::subject_value(MetadataRecordSet::PolicyDecisionCurrentById, false)]
202-
fn test_migration_rewrites_legacy_policy_subjects(#[case] record_set: MetadataRecordSet, #[case] key: bool) {
201+
#[test]
202+
fn test_migration_rewrites_legacy_policy_subject_keys() {
203203
let subject = json!({
204204
"repository": "repo",
205205
"project": "resource",
@@ -208,34 +208,27 @@ fn test_migration_rewrites_legacy_policy_subjects(#[case] record_set: MetadataRe
208208
"source": "upstream",
209209
"action": "serve"
210210
});
211-
let record = if key {
212-
MetadataRecord {
211+
let migrated = rewritten(
212+
MetadataRecordSet::PolicyDecisionCurrent,
213+
&MetadataRecord {
213214
key: subject.to_string(),
214-
value: b"decision".to_vec(),
215-
}
216-
} else {
217-
MetadataRecord {
218-
key: "decision".to_owned(),
219-
value: serde_json::to_vec(&subject).unwrap(),
220-
}
221-
};
222-
let migrated = rewritten(record_set, &record);
223-
let migrated_subject = if key {
224-
serde_json::from_str(&migrated.key).unwrap()
225-
} else {
226-
value(&migrated)
227-
};
215+
value: b"pd_0000000000000001".to_vec(),
216+
},
217+
);
228218

229219
assert_eq!(
230-
migrated_subject,
231-
json!({
232-
"repository": "repo",
233-
"resource": "resource",
234-
"group": "group",
235-
"artifact": "artifact",
236-
"source": "upstream",
237-
"action": "serve"
238-
})
220+
(serde_json::from_str::<Value>(&migrated.key).unwrap(), migrated.value),
221+
(
222+
json!({
223+
"repository": "repo",
224+
"resource": "resource",
225+
"group": "group",
226+
"artifact": "artifact",
227+
"source": "upstream",
228+
"action": "serve"
229+
}),
230+
b"pd_0000000000000001".to_vec()
231+
)
239232
);
240233
}
241234

@@ -365,6 +358,25 @@ fn test_migration_rewrites_legacy_daily_analytics() {
365358
"next_eligible_at_unix": 10
366359
})
367360
)]
361+
#[case::current_policy_by_id(
362+
MetadataRecordSet::PolicyDecisionCurrentById,
363+
"decision",
364+
json!({
365+
"id": "00000000-0000-0000-0000-000000000002",
366+
"repository": "repo",
367+
"resource": "resource",
368+
"group": "group",
369+
"artifact": "artifact",
370+
"source": "upstream",
371+
"action": "serve",
372+
"state": "allow",
373+
"rule": "rule",
374+
"reason": "reason",
375+
"evaluated_at_unix": 9,
376+
"input_generation": {"repository": 1, "catalog": 2, "policy": 3},
377+
"next_eligible_at_unix": 10
378+
})
379+
)]
368380
#[case::current_reads(MetadataRecordSet::Analytics, "reads", json!({"artifacts": []}))]
369381
#[case::current_daily(MetadataRecordSet::Analytics, "daily_usage", json!({"schema": 1, "buckets": []}))]
370382
#[case::unknown_analytics(MetadataRecordSet::Analytics, "other", json!({}))]
@@ -377,10 +389,8 @@ fn test_migration_leaves_current_or_unowned_records_unchanged(
377389
assert_eq!(PypiPlugin.rewrite(record_set, &record(key, contents)), Ok(None));
378390
}
379391

380-
#[rstest::rstest]
381-
#[case::subject_key(MetadataRecordSet::PolicyDecisionCurrent, true)]
382-
#[case::subject_value(MetadataRecordSet::PolicyDecisionCurrentById, false)]
383-
fn test_migration_leaves_current_policy_subjects_unchanged(#[case] record_set: MetadataRecordSet, #[case] key: bool) {
392+
#[test]
393+
fn test_migration_leaves_current_policy_subject_keys_unchanged() {
384394
let subject = json!({
385395
"repository": "repo",
386396
"resource": "resource",
@@ -389,38 +399,31 @@ fn test_migration_leaves_current_policy_subjects_unchanged(#[case] record_set: M
389399
"source": "upstream",
390400
"action": "serve"
391401
});
392-
let record = if key {
393-
MetadataRecord {
394-
key: subject.to_string(),
395-
value: b"decision".to_vec(),
396-
}
397-
} else {
398-
MetadataRecord {
399-
key: "decision".to_owned(),
400-
value: serde_json::to_vec(&subject).unwrap(),
401-
}
402-
};
403402

404-
assert_eq!(PypiPlugin.rewrite(record_set, &record), Ok(None));
403+
assert_eq!(
404+
PypiPlugin.rewrite(
405+
MetadataRecordSet::PolicyDecisionCurrent,
406+
&MetadataRecord {
407+
key: subject.to_string(),
408+
value: b"pd_0000000000000001".to_vec(),
409+
}
410+
),
411+
Ok(None)
412+
);
405413
}
406414

407-
#[rstest::rstest]
408-
#[case::subject_key(MetadataRecordSet::PolicyDecisionCurrent, true)]
409-
#[case::subject_value(MetadataRecordSet::PolicyDecisionCurrentById, false)]
410-
fn test_migration_leaves_malformed_policy_subjects_unchanged(#[case] record_set: MetadataRecordSet, #[case] key: bool) {
411-
let record = if key {
412-
MetadataRecord {
413-
key: "malformed".to_owned(),
414-
value: b"decision".to_vec(),
415-
}
416-
} else {
417-
MetadataRecord {
418-
key: "decision".to_owned(),
419-
value: b"malformed".to_vec(),
420-
}
421-
};
422-
423-
assert_eq!(PypiPlugin.rewrite(record_set, &record), Ok(None));
415+
#[test]
416+
fn test_migration_leaves_malformed_policy_subject_keys_unchanged() {
417+
assert_eq!(
418+
PypiPlugin.rewrite(
419+
MetadataRecordSet::PolicyDecisionCurrent,
420+
&MetadataRecord {
421+
key: "malformed".to_owned(),
422+
value: b"pd_0000000000000001".to_vec(),
423+
}
424+
),
425+
Ok(None)
426+
);
424427
}
425428

426429
#[test]
@@ -518,6 +521,28 @@ fn write_policy_fixture(path: &Path) {
518521
}),
519522
)],
520523
);
524+
write_bytes(
525+
path,
526+
"policy_decision_current_id",
527+
&[(
528+
"decision",
529+
json!({
530+
"id": "00000000-0000-0000-0000-000000000002",
531+
"repository": "repo",
532+
"project": "resource",
533+
"version": "group",
534+
"filename": "artifact",
535+
"source": "upstream",
536+
"action": "serve",
537+
"state": "allow",
538+
"rule": null,
539+
"reason": null,
540+
"evaluated_at_unix": 15,
541+
"input_generation": {"repository": 1, "catalog": 2, "policy": 3},
542+
"next_eligible_at_unix": null
543+
}),
544+
)],
545+
);
521546
let legacy_subject = json!({
522547
"repository": "repo",
523548
"project": "resource",
@@ -531,11 +556,6 @@ fn write_policy_fixture(path: &Path) {
531556
"policy_decision_current",
532557
&[(legacy_subject.to_string(), "decision".to_owned())],
533558
);
534-
write_text(
535-
path,
536-
"policy_decision_current_id",
537-
&[("decision".to_owned(), legacy_subject.to_string())],
538-
);
539559
}
540560

541561
fn write_analytics_fixture(path: &Path) {
@@ -688,7 +708,7 @@ fn assert_migrated_policy(path: &Path) {
688708
)
689709
);
690710
assert_eq!(
691-
value_bytes(read_text(path, "policy_decision_current_id")["decision"].as_bytes())["resource"],
711+
value_bytes(&read_bytes(path, "policy_decision_current_id")["decision"])["resource"],
692712
"resource"
693713
);
694714
}

crates/peryx-ecosystem-pypi/tests/unit/shadow/coverage_tests.rs

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -380,7 +380,7 @@ impl MetadataMigration for UnreadableDecision {
380380
}
381381

382382
fn record_sets(&self) -> &[MetadataRecordSet] {
383-
&[MetadataRecordSet::PolicyDecisionHistory]
383+
&[MetadataRecordSet::PolicyDecisionCurrentById]
384384
}
385385

386386
fn legacy_sources(&self) -> &[LegacyMetadataSource] {

crates/peryx-storage/src/meta/migration.rs

Lines changed: 8 additions & 20 deletions
Original file line numberDiff line numberDiff line change
@@ -358,7 +358,7 @@ fn collect_target(
358358
collect_text(txn, POLICY_DECISION_CURRENT, cursor, Some(created_keys))
359359
}
360360
MetadataRecordSet::PolicyDecisionCurrentById => {
361-
collect_text(txn, POLICY_DECISION_CURRENT_ID, cursor, Some(created_keys))
361+
collect_bytes(txn, POLICY_DECISION_CURRENT_ID, cursor, Some(created_keys))
362362
}
363363
MetadataRecordSet::Analytics => collect_bytes(txn, ANALYTICS, cursor, Some(created_keys)),
364364
}
@@ -443,14 +443,9 @@ fn rewrite_target(
443443
MetadataRecordSet::PolicyDecisionCurrent => {
444444
rewrite_text(txn, POLICY_DECISION_CURRENT, original, rewritten, migration, record_set)
445445
}
446-
MetadataRecordSet::PolicyDecisionCurrentById => rewrite_text(
447-
txn,
448-
POLICY_DECISION_CURRENT_ID,
449-
original,
450-
rewritten,
451-
migration,
452-
record_set,
453-
),
446+
MetadataRecordSet::PolicyDecisionCurrentById => {
447+
rewrite_bytes(txn, POLICY_DECISION_CURRENT_ID, original, rewritten).map_err(Into::into)
448+
}
454449
MetadataRecordSet::Analytics => rewrite_bytes(txn, ANALYTICS, original, rewritten).map_err(Into::into),
455450
}
456451
}
@@ -471,14 +466,9 @@ fn insert_target(
471466
MetadataRecordSet::PolicyDecisionCurrent => {
472467
insert_text(txn, POLICY_DECISION_CURRENT, original, rewritten, migration, record_set)
473468
}
474-
MetadataRecordSet::PolicyDecisionCurrentById => insert_text(
475-
txn,
476-
POLICY_DECISION_CURRENT_ID,
477-
original,
478-
rewritten,
479-
migration,
480-
record_set,
481-
),
469+
MetadataRecordSet::PolicyDecisionCurrentById => {
470+
insert_bytes(txn, POLICY_DECISION_CURRENT_ID, rewritten).map_err(Into::into)
471+
}
482472
MetadataRecordSet::Analytics => insert_bytes(txn, ANALYTICS, rewritten).map_err(Into::into),
483473
}
484474
}
@@ -623,9 +613,7 @@ const fn target_name(record_set: MetadataRecordSet) -> &'static str {
623613

624614
const fn target_value_kind(record_set: MetadataRecordSet) -> MetadataValueKind {
625615
match record_set {
626-
MetadataRecordSet::PolicyDecisionCurrent | MetadataRecordSet::PolicyDecisionCurrentById => {
627-
MetadataValueKind::Text
628-
}
616+
MetadataRecordSet::PolicyDecisionCurrent => MetadataValueKind::Text,
629617
_ => MetadataValueKind::Bytes,
630618
}
631619
}

crates/peryx-storage/src/meta/mod.rs

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -131,9 +131,13 @@ const INGRESS_INTENT_ORDER: TableDefinition<u64, &str> = TableDefinition::new("i
131131
const INGRESS_INTENT_SEQ: TableDefinition<&str, u64> = TableDefinition::new("ingress_intent_seq");
132132
const INGRESS_SEQ_KEY: &str = "next";
133133
const RECONCILE_BACKLOG: TableDefinition<&str, &[u8]> = TableDefinition::new("reconcile_backlog");
134+
/// Bounded audit log of every evaluation, oldest rows evicted once it crosses its limit.
134135
const POLICY_DECISION: TableDefinition<&str, &[u8]> = TableDefinition::new("policy_decision");
136+
/// Subject to the identifier of the decision that currently holds for it.
135137
const POLICY_DECISION_CURRENT: TableDefinition<&str, &str> = TableDefinition::new("policy_decision_current");
136-
const POLICY_DECISION_CURRENT_ID: TableDefinition<&str, &str> = TableDefinition::new("policy_decision_current_id");
138+
/// Current decisions in evaluation order, stored apart from the audit log so that evicting audit
139+
/// rows cannot drop live subject state.
140+
const POLICY_DECISION_CURRENT_ID: TableDefinition<&str, &[u8]> = TableDefinition::new("policy_decision_current_id");
137141
const POLICY_INPUT_GENERATION: TableDefinition<&str, &[u8]> = TableDefinition::new("policy_input_generation");
138142
const DERIVED_VIEW_FRONTIER: TableDefinition<&str, u64> = TableDefinition::new("derived_view_frontier");
139143
const QUOTA_USAGE: TableDefinition<&str, &[u8]> = TableDefinition::new("quota_usage");

0 commit comments

Comments
 (0)