michal/tit
Browse tree · Show commit · Download archive
Diff
33b8025ee5c6 → 70b42e946c6f
Cargo.lock
Mode 100644 → 100644; object 4a555ba65cff → 86912b552c6b
@@ -3579,6 +3579,7 @@ "rusqlite", "russh", "serde", + "serde_json", "sha2", "ssh-key", "tempfile",
Cargo.toml
Mode 100644 → 100644; object 7112bfe27085 → 7f4dd1a9d471
@@ -25,6 +25,7 @@
rusqlite = { version = "0.40", default-features = false, features = ["backup", "bundled"] }
russh = { version = "0.62", default-features = false, features = ["ring"] }
serde = { version = "1.0", features = ["derive"] }
+serde_json = "1.0"
sha2 = { version = "0.11", default-features = false }
ssh-key = { version = "=0.7.0-rc.11", default-features = false, features = ["ed25519", "p256", "std"] }
thiserror = "2.0"
README.md
Mode 100644 → 100644; object 380f3cd45a46 → 9d37740daab8
@@ -292,3 +292,17 @@ This command tests one Web and SSH identity, account recovery, key revocation, sessions, repository roles, private route isolation, push policy, audit history, and repository creation from the Web UI and SSH. + +## Milestone 4.1 gate + +Run the repository event service gate: + +```text +./scripts/check-m4-1 +``` + +This command tests event migration, random public IDs, repository sequences, +versioned JSON payloads, atomic metadata and event writes, Git operation +recovery, feed parsing, and sequence pagination. Read the +[repository event service architectural decision record](docs/adr/0016-repository-event-service.md) +for the event type and payload contracts.
docs/adr/0007-public-repository-feeds.md
Mode 100644 → 100644; object 1f4909617b11 → 9c6d1984b5b1
@@ -1,6 +1,6 @@ # Architectural decision record 0007: public repository feeds -Status: Accepted +Status: Superseded by architectural decision record 0016 Date: 2026-07-23
docs/adr/0016-repository-event-service.md
Mode → 100644; object → b450b606f6d0
@@ -1,0 +1,97 @@ +# Architectural decision record 0016: Repository event service + +Status: Accepted + +Date: 2026-07-22 + +## Context + +Repository changes need one durable event source for feeds and subsequent +subscriptions. A database row ID is not a public identifier. A global row order +also does not give a repository-local cursor. + +Each event payload needs a version before issue and pull-request events add new +schemas. An event must not exist without its metadata change. A metadata change +must not exist without its event. + +## Decision + +Keep `repository_event` as the one canonical repository event table. Do not add +a parallel event history. Give each event these identifiers: + +- `event_id` is a random 128-bit lowercase hexadecimal public ID. +- `sequence` is an integer that increases by one inside each repository. +- `id` stays as a private SQLite row key. + +Allocate the next sequence in the same immediate transaction that inserts the +event. A unique constraint on `(repository_id, sequence)` prevents reuse. Read +event pages through the `(repository_id, sequence DESC)` index. Use the +repository sequence as the `before` cursor. + +The stable version 1 event types are `repository-created`, +`repository-imported`, `push`, `ref-created`, `ref-updated`, `ref-deleted`, +`tag-created`, `tag-updated`, and `tag-deleted`. + +Store `payload_version` as an explicit schema value and store `version` in each +JSON payload. The schema accepts only a JSON object whose inner version equals +the column value. Limit a payload to 1 MiB. Version 1 repository payloads have +`owner`, `repository`, and `object_format`. Push payloads have `operation_id`. +Ref and tag payloads have `name_hex`, `old_target`, and `new_target`. Hexadecimal +ref names preserve Git ref bytes without a text conversion. + +Use one insertion helper for repository creation, repository import, initial +refs and tags, a completed push, and changed refs and tags. The helper runs in +the transaction that owns the related metadata mutation. If event insertion +fails, roll back the metadata change. A Git operation stays incomplete if its +completion events cannot be inserted. + +Migration 011 assigns a random public ID and a repository sequence to each old +event. It creates a version 1 payload from the old typed columns. The old row ID +does not become the public ID. + +Atom and RSS entry IDs now use this form: + +```text +urn:tit:event:EVENT_ID +``` + +The value stays stable after a restart, repository rename, or later event +insert. Existing feed pages use the repository sequence for pagination. + +## Failure and threat cases + +Database constraints reject an invalid public ID, repeated sequence, unknown +event type, unsupported payload version, malformed JSON, non-object JSON, +missing inner version, mismatched version, oversized payload, invalid ref +shape, and invalid Git operation source. + +The event payload is application output. Repository and account names pass the +domain validators. Ref names use hexadecimal text. Object IDs and operation IDs +use their validated canonical values. Event payload generation does not accept +an arbitrary JSON string from an HTTP or SSH request. + +Git refs and SQLite still use the durable Git operation intent. The server +inserts push events only when it marks the intent complete. Restart recovery +uses the same completion transaction, so a feed cannot announce a push whose +refs are not reachable. + +## Evidence + +Migration tests upgrade each historical schema, including schema version 10. +They also kill a migration before and after commit. Storage tests verify random +public IDs, consecutive repository sequences, versioned JSON schemas, indexed +pagination, stable IDs after reopen, duplicate-sequence rejection, payload +rejection, and backfill. + +Injected trigger failures prove that repository creation rolls back when its +event fails. They also prove that a Git operation stays incomplete and publishes +no events when its completion event fails. Production feed tests parse Atom and +RSS and follow sequence pagination. Git crash tests continue to stop the process +at each cross-store boundary. + +## Consequences + +Issue and pull-request milestones can add event types and payload versions to +this service. Consumers use a public event ID for identity and a repository +sequence for order. The private SQLite row key can change without changing +either contract.
scripts/check-m4-1
Mode → 100755; object → c1fd5d8567df
@@ -1,0 +1,7 @@ +#!/bin/sh +set -eu + +./scripts/check +cargo test --locked --release --test sqlite +cargo test --locked --release --test public_routes browses_and_clones_public_repositories_for_both_hash_formats +cargo test --locked --release --test git_push_ssh reconciles_process_termination_at_each_cross_store_boundary
src/feed.rs
Mode 100644 → 100644; object ca31a497f93a → cb163878fc69
@@ -50,7 +50,7 @@
}
for event in self.events {
output.push_str("<entry>\n");
- element(&mut output, "id", &event_id(&self.repository.id, event.id))?;
+ element(&mut output, "id", &event_id(&event.event_id))?;
element(&mut output, "title", &event_title(event))?;
element(&mut output, "updated", &atom_date(event.created_at)?)?;
empty_link(&mut output, "alternate", &repository_url)?;
@@ -92,7 +92,7 @@
element(&mut output, "title", &event_title(event))?;
element(&mut output, "link", &repository_url)?;
write!(output, "<guid isPermaLink=\"false\">")?;
- escape_xml(&event_id(&self.repository.id, event.id), &mut output)?;
+ escape_xml(&event_id(&event.event_id), &mut output)?;
output.push_str("</guid>\n");
element(&mut output, "pubDate", &rss_date(event.created_at)?)?;
element(&mut output, "description", &event_description(event))?;
@@ -123,8 +123,8 @@
}
}
-fn event_id(repository_id: &str, event_id: i64) -> String {
- format!("urn:tit:event:{repository_id}:{event_id}")
+fn event_id(event_id: &str) -> String {
+ format!("urn:tit:event:{event_id}")
}
fn event_title(event: &RepositoryEventRecord) -> String {
src/http/public.rs
Mode 100644 → 100644; object 748c7c6ff631 → dacddf00f825
@@ -333,7 +333,7 @@
let has_next = events.len() > PAGE_SIZE;
events.truncate(PAGE_SIZE);
let next_before = has_next
- .then(|| events.last().map(|event| event.id))
+ .then(|| events.last().map(|event| event.sequence))
.flatten();
let name = match format {
FeedFormat::Atom => "atom.xml",
src/store/event.rs
Mode → 100644; object → 001aea2a97b6
@@ -1,0 +1,106 @@
+use serde_json::json;
+
+pub(super) const PAYLOAD_VERSION: i64 = 1;
+
+#[derive(Clone, Copy)]
+pub(super) enum EventKind {
+ RepositoryCreated,
+ RepositoryImported,
+ Push,
+ RefCreated,
+ RefUpdated,
+ RefDeleted,
+ TagCreated,
+ TagUpdated,
+ TagDeleted,
+}
+
+impl EventKind {
+ pub(super) fn as_str(self) -> &'static str {
+ match self {
+ Self::RepositoryCreated => "repository-created",
+ Self::RepositoryImported => "repository-imported",
+ Self::Push => "push",
+ Self::RefCreated => "ref-created",
+ Self::RefUpdated => "ref-updated",
+ Self::RefDeleted => "ref-deleted",
+ Self::TagCreated => "tag-created",
+ Self::TagUpdated => "tag-updated",
+ Self::TagDeleted => "tag-deleted",
+ }
+ }
+}
+
+pub(super) struct VersionedEvent {
+ pub(super) kind: EventKind,
+ pub(super) payload: String,
+}
+
+pub(super) fn repository(
+ kind: EventKind,
+ owner: &str,
+ repository: &str,
+ object_format: &str,
+) -> VersionedEvent {
+ debug_assert!(matches!(
+ kind,
+ EventKind::RepositoryCreated | EventKind::RepositoryImported
+ ));
+ VersionedEvent {
+ kind,
+ payload: json!({
+ "version": PAYLOAD_VERSION,
+ "owner": owner,
+ "repository": repository,
+ "object_format": object_format,
+ })
+ .to_string(),
+ }
+}
+
+pub(super) fn push(operation_id: &str) -> VersionedEvent {
+ VersionedEvent {
+ kind: EventKind::Push,
+ payload: json!({
+ "version": PAYLOAD_VERSION,
+ "operation_id": operation_id,
+ })
+ .to_string(),
+ }
+}
+
+pub(super) fn reference(
+ kind: EventKind,
+ name: &[u8],
+ old_target: Option<&str>,
+ new_target: Option<&str>,
+) -> VersionedEvent {
+ debug_assert!(matches!(
+ kind,
+ EventKind::RefCreated
+ | EventKind::RefUpdated
+ | EventKind::RefDeleted
+ | EventKind::TagCreated
+ | EventKind::TagUpdated
+ | EventKind::TagDeleted
+ ));
+ VersionedEvent {
+ kind,
+ payload: json!({
+ "version": PAYLOAD_VERSION,
+ "name_hex": encode_hex(name),
+ "old_target": old_target,
+ "new_target": new_target,
+ })
+ .to_string(),
+ }
+}
+
+fn encode_hex(bytes: &[u8]) -> String {
+ let mut encoded = String::with_capacity(bytes.len().saturating_mul(2));
+ for byte in bytes {
+ use std::fmt::Write;
+ write!(encoded, "{byte:02x}").expect("a string write cannot fail");
+ }
+ encoded
+}
src/store/migrations/011_domain_events.sql
Mode → 100644; object → 8250ab8b1229
@@ -1,0 +1,98 @@
+DROP INDEX repository_event_feed;
+ALTER TABLE repository_event RENAME TO repository_event_v10;
+
+CREATE TABLE repository_event (
+ id INTEGER PRIMARY KEY,
+ event_id TEXT NOT NULL UNIQUE
+ CHECK (
+ length(event_id) = 32
+ AND event_id = lower(event_id)
+ AND event_id NOT GLOB '*[^0-9a-f]*'
+ ),
+ repository_id TEXT NOT NULL
+ REFERENCES repository (id) ON DELETE RESTRICT,
+ sequence INTEGER NOT NULL CHECK (sequence >= 1),
+ source_intent_id TEXT
+ REFERENCES git_operation_intent (id) ON DELETE RESTRICT,
+ source_ordinal INTEGER CHECK (source_ordinal IS NULL OR source_ordinal >= 0),
+ kind TEXT NOT NULL
+ CHECK (kind IN (
+ 'repository-created', 'repository-imported', 'push',
+ 'ref-created', 'ref-updated', 'ref-deleted',
+ 'tag-created', 'tag-updated', 'tag-deleted'
+ )),
+ actor TEXT NOT NULL CHECK (length(actor) BETWEEN 1 AND 256),
+ ref_name BLOB,
+ old_target TEXT,
+ new_target TEXT,
+ payload_version INTEGER NOT NULL CHECK (payload_version = 1),
+ payload TEXT NOT NULL
+ CHECK (
+ length(payload) BETWEEN 1 AND 1048576
+ AND CASE WHEN json_valid(payload) THEN
+ json_type(payload) = 'object'
+ AND coalesce(json_extract(payload, '$.version') = payload_version, 0)
+ ELSE 0 END
+ ),
+ created_at INTEGER NOT NULL CHECK (created_at >= 0),
+ UNIQUE (repository_id, sequence),
+ UNIQUE (source_intent_id, source_ordinal),
+ CHECK (
+ (source_intent_id IS NULL AND source_ordinal IS NULL)
+ OR (source_intent_id IS NOT NULL AND source_ordinal IS NOT NULL)
+ ),
+ CHECK (
+ (kind IN ('repository-created', 'repository-imported', 'push')
+ AND ref_name IS NULL AND old_target IS NULL AND new_target IS NULL)
+ OR
+ (kind LIKE 'ref-%' OR kind LIKE 'tag-%')
+ AND ref_name IS NOT NULL
+ AND (old_target IS NOT NULL OR new_target IS NOT NULL)
+ )
+) STRICT;
+
+INSERT INTO repository_event
+ (id, event_id, repository_id, sequence, source_intent_id, source_ordinal,
+ kind, actor, ref_name, old_target, new_target, payload_version, payload,
+ created_at)
+SELECT
+ legacy.id,
+ lower(hex(randomblob(16))),
+ legacy.repository_id,
+ row_number() OVER (PARTITION BY legacy.repository_id ORDER BY legacy.id),
+ legacy.source_intent_id,
+ legacy.source_ordinal,
+ legacy.kind,
+ legacy.actor,
+ legacy.ref_name,
+ legacy.old_target,
+ legacy.new_target,
+ 1,
+ CASE
+ WHEN legacy.kind IN ('repository-created', 'repository-imported') THEN
+ json_object(
+ 'version', 1,
+ 'owner', account.username,
+ 'repository', repository.slug,
+ 'object_format', repository.object_format
+ )
+ WHEN legacy.kind = 'push' THEN
+ json_object('version', 1, 'operation_id', legacy.source_intent_id)
+ ELSE
+ json_object(
+ 'version', 1,
+ 'name_hex', lower(hex(legacy.ref_name)),
+ 'old_target', legacy.old_target,
+ 'new_target', legacy.new_target
+ )
+ END,
+ legacy.created_at
+FROM repository_event_v10 AS legacy
+JOIN repository ON repository.id = legacy.repository_id
+JOIN account ON account.id = repository.owner_account_id
+ORDER BY legacy.id;
+
+DROP TABLE repository_event_v10;
+
+CREATE INDEX repository_event_feed
+ON repository_event (repository_id, sequence DESC);
src/store/mod.rs
Mode 100644 → 100644; object 1beabf69979e → b05170aa762a
@@ -7,9 +7,11 @@
use rusqlite::{Connection, OptionalExtension, TransactionBehavior};
use thiserror::Error;
+mod event;
+
const BUSY_TIMEOUT: Duration = Duration::from_secs(5);
const BUSY_TIMEOUT_MILLISECONDS: i64 = 5_000;
-const SCHEMA_VERSION: i64 = 10;
+const SCHEMA_VERSION: i64 = 11;
#[allow(
dead_code,
reason = "the integration test imports this module without the CLI operation"
@@ -19,7 +21,7 @@
dead_code,
reason = "M1A proves migrations before the M2 server calls them"
)]
-const MIGRATIONS: [&str; 10] = [
+const MIGRATIONS: [&str; 11] = [
include_str!("migrations/001_initial.sql"),
include_str!("migrations/002_state.sql"),
include_str!("migrations/003_git_intents.sql"),
@@ -30,6 +32,7 @@
include_str!("migrations/008_web_sessions.sql"),
include_str!("migrations/009_repository_authorization.sql"),
include_str!("migrations/010_audit_history.sql"),
+ include_str!("migrations/011_domain_events.sql"),
];
#[allow(
@@ -936,37 +939,49 @@
);
match result {
Ok(1) => {
- transaction.execute(
- "INSERT INTO repository_event
- (repository_id, kind, actor, created_at)
- VALUES (?1, ?2, ?3, ?4)",
- rusqlite::params![
- repository.id,
- repository.origin.event_kind(),
- repository.owner,
- repository.created_at,
- ],
+ let event = event::repository(
+ repository.origin.event_kind(),
+ repository.owner,
+ repository.slug,
+ repository.object_format,
+ );
+ insert_domain_event(
+ &transaction,
+ &NewDomainEvent {
+ repository_id: repository.id,
+ source_intent_id: None,
+ source_ordinal: None,
+ event: &event,
+ actor: repository.owner,
+ ref_name: None,
+ old_target: None,
+ new_target: None,
+ created_at: repository.created_at,
+ },
)?;
for reference in repository.initial_references {
let kind = if reference.name.starts_with(b"refs/tags/") {
- "tag-created"
+ event::EventKind::TagCreated
} else if reference.name.starts_with(b"refs/heads/") {
- "ref-created"
+ event::EventKind::RefCreated
} else {
return Err(StoreError::EventPayload);
};
- transaction.execute(
- "INSERT INTO repository_event
- (repository_id, kind, actor, ref_name, new_target, created_at)
- VALUES (?1, ?2, ?3, ?4, ?5, ?6)",
- rusqlite::params![
- repository.id,
- kind,
- repository.owner,
- reference.name,
- reference.target,
- repository.created_at,
- ],
+ let event =
+ event::reference(kind, &reference.name, None, Some(&reference.target));
+ insert_domain_event(
+ &transaction,
+ &NewDomainEvent {
+ repository_id: repository.id,
+ source_intent_id: None,
+ source_ordinal: None,
+ event: &event,
+ actor: repository.owner,
+ ref_name: Some(&reference.name),
+ old_target: None,
+ new_target: Some(&reference.target),
+ created_at: repository.created_at,
+ },
)?;
}
let target = format!("{}/{}", repository.owner, repository.slug);
@@ -1487,22 +1502,26 @@
) -> Result<(RepositoryRecord, Vec<RepositoryEventRecord>), StoreError> {
let limit = i64::try_from(limit).map_err(|_| StoreError::EventLimit)?;
let mut statement = self.connection.prepare(
- "SELECT id, kind, actor, ref_name, old_target, new_target, created_at
+ "SELECT event_id, sequence, kind, actor, ref_name, old_target, new_target,
+ payload_version, payload, created_at
FROM repository_event
- WHERE repository_id = ?1 AND (?2 IS NULL OR id < ?2)
- ORDER BY id DESC
+ WHERE repository_id = ?1 AND (?2 IS NULL OR sequence < ?2)
+ ORDER BY sequence DESC
LIMIT ?3",
)?;
let events = statement
.query_map(rusqlite::params![repository.id, before, limit], |row| {
Ok(RepositoryEventRecord {
- id: row.get(0)?,
- kind: row.get(1)?,
- actor: row.get(2)?,
- ref_name: row.get(3)?,
- old_target: row.get(4)?,
- new_target: row.get(5)?,
- created_at: row.get(6)?,
+ event_id: row.get(0)?,
+ sequence: row.get(1)?,
+ kind: row.get(2)?,
+ actor: row.get(3)?,
+ ref_name: row.get(4)?,
+ old_target: row.get(5)?,
+ new_target: row.get(6)?,
+ payload_version: row.get(7)?,
+ payload: row.get(8)?,
+ created_at: row.get(9)?,
})
})?
.collect::<Result<Vec<_>, _>>()?;
@@ -1616,10 +1635,10 @@
}
impl RepositoryOrigin {
- fn event_kind(self) -> &'static str {
+ fn event_kind(self) -> event::EventKind {
match self {
- Self::Created => "repository-created",
- Self::Imported => "repository-imported",
+ Self::Created => event::EventKind::RepositoryCreated,
+ Self::Imported => event::EventKind::RepositoryImported,
}
}
@@ -1672,12 +1691,15 @@
reason = "some integration tests import storage without public event pages"
)]
pub(crate) struct RepositoryEventRecord {
- pub(crate) id: i64,
+ pub(crate) event_id: String,
+ pub(crate) sequence: i64,
pub(crate) kind: String,
pub(crate) actor: String,
pub(crate) ref_name: Option<Vec<u8>>,
pub(crate) old_target: Option<String>,
pub(crate) new_target: Option<String>,
+ pub(crate) payload_version: i64,
+ pub(crate) payload: String,
pub(crate) created_at: i64,
}
@@ -1780,6 +1802,51 @@
Ok(())
}
+struct NewDomainEvent<'a> {
+ repository_id: &'a str,
+ source_intent_id: Option<&'a str>,
+ source_ordinal: Option<i64>,
+ event: &'a event::VersionedEvent,
+ actor: &'a str,
+ ref_name: Option<&'a [u8]>,
+ old_target: Option<&'a str>,
+ new_target: Option<&'a str>,
+ created_at: i64,
+}
+
+fn insert_domain_event(
+ transaction: &rusqlite::Transaction<'_>,
+ event: &NewDomainEvent<'_>,
+) -> Result<(), StoreError> {
+ transaction.execute(
+ "INSERT INTO repository_event
+ (event_id, repository_id, sequence, source_intent_id, source_ordinal,
+ kind, actor, ref_name, old_target, new_target, payload_version, payload,
+ created_at)
+ VALUES (
+ lower(hex(randomblob(16))),
+ ?1,
+ (SELECT COALESCE(MAX(sequence), 0) + 1
+ FROM repository_event WHERE repository_id = ?1),
+ ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11
+ )",
+ rusqlite::params![
+ event.repository_id,
+ event.source_intent_id,
+ event.source_ordinal,
+ event.event.kind.as_str(),
+ event.actor,
+ event.ref_name,
+ event.old_target,
+ event.new_target,
+ event::PAYLOAD_VERSION,
+ event.event.payload,
+ event.created_at,
+ ],
+ )?;
+ Ok(())
+}
+
fn end_sessions(
transaction: &rusqlite::Transaction<'_>,
account_id: i64,
@@ -1821,51 +1888,72 @@
if initial.len() != proposed.len() {
return Err(StoreError::EventPayload);
}
- transaction.execute(
- "INSERT INTO repository_event
- (repository_id, source_intent_id, source_ordinal, kind, actor, created_at)
- VALUES (?1, ?2, 0, 'push', ?3, ?4)",
- rusqlite::params![repository_id, intent_id, actor, created_at],
+ let push = event::push(intent_id);
+ insert_domain_event(
+ transaction,
+ &NewDomainEvent {
+ repository_id,
+ source_intent_id: Some(intent_id),
+ source_ordinal: Some(0),
+ event: &push,
+ actor,
+ ref_name: None,
+ old_target: None,
+ new_target: None,
+ created_at,
+ },
)?;
for (index, ((old, old_name), (new, new_name))) in initial.into_iter().zip(proposed).enumerate()
{
if old_name != new_name || (is_null_id(&old) && is_null_id(&new)) {
return Err(StoreError::EventPayload);
}
- let prefix = if old_name.starts_with(b"refs/tags/") {
- "tag"
+ let tag = if old_name.starts_with(b"refs/tags/") {
+ true
} else if old_name.starts_with(b"refs/heads/") {
- "ref"
+ false
} else {
return Err(StoreError::EventPayload);
};
- let action = if is_null_id(&old) {
- "created"
+ let kind = if is_null_id(&old) {
+ if tag {
+ event::EventKind::TagCreated
+ } else {
+ event::EventKind::RefCreated
+ }
} else if is_null_id(&new) {
- "deleted"
+ if tag {
+ event::EventKind::TagDeleted
+ } else {
+ event::EventKind::RefDeleted
+ }
+ } else if tag {
+ event::EventKind::TagUpdated
} else {
- "updated"
+ event::EventKind::RefUpdated
};
- let kind = format!("{prefix}-{action}");
let old_target = (!is_null_id(&old)).then_some(old);
let new_target = (!is_null_id(&new)).then_some(new);
let ordinal = i64::try_from(index + 1).map_err(|_| StoreError::EventPayload)?;
- transaction.execute(
- "INSERT INTO repository_event
- (repository_id, source_intent_id, source_ordinal, kind, actor,
- ref_name, old_target, new_target, created_at)
- VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9)",
- rusqlite::params![
+ let event = event::reference(
+ kind,
+ &old_name,
+ old_target.as_deref(),
+ new_target.as_deref(),
+ );
+ insert_domain_event(
+ transaction,
+ &NewDomainEvent {
repository_id,
- intent_id,
- ordinal,
- kind,
+ source_intent_id: Some(intent_id),
+ source_ordinal: Some(ordinal),
+ event: &event,
actor,
- old_name,
- old_target,
- new_target,
+ ref_name: Some(&old_name),
+ old_target: old_target.as_deref(),
+ new_target: new_target.as_deref(),
created_at,
- ],
+ },
)?;
}
insert_audit_event(
tests/cli.rs
Mode 100644 → 100644; object d9ecb27c142a → 1255ad2d777f
@@ -19,7 +19,8 @@
include_str!("../src/store/migrations/008_web_sessions.sql"),
include_str!("../src/store/migrations/009_repository_authorization.sql"),
include_str!("../src/store/migrations/010_audit_history.sql"),
- "PRAGMA user_version = 10;\n",
+ include_str!("../src/store/migrations/011_domain_events.sql"),
+ "PRAGMA user_version = 11;\n",
);
#[test]
tests/public_routes.rs
Mode 100644 → 100644; object e33ea94466da → 052d69b9b0b9
@@ -173,8 +173,15 @@
.connection()
.execute(
"INSERT INTO repository_event
- (repository_id, kind, actor, created_at)
- VALUES (?1, 'push', ?2, ?3)",
+ (event_id, repository_id, sequence, kind, actor, payload_version,
+ payload, created_at)
+ VALUES (
+ lower(hex(randomblob(16))),
+ ?1,
+ (SELECT COALESCE(MAX(sequence), 0) + 1
+ FROM repository_event WHERE repository_id = ?1),
+ 'push', ?2, 1, '{\"version\":1}', ?3
+ )",
rusqlite::params![fixture.repository_id, actor, timestamp],
)
.expect("insert a page fixture event");
tests/sqlite.rs
Mode 100644 → 100644; object 4a33b211e1f4 → 890c6879b579
@@ -33,6 +33,13 @@
include_str!("../src/store/migrations/009_repository_authorization.sql"),
"PRAGMA user_version = 9;\n",
);
+const V10_FIXTURE: &str = concat!(
+ include_str!("fixtures/sqlite/v7.sql"),
+ include_str!("../src/store/migrations/008_web_sessions.sql"),
+ include_str!("../src/store/migrations/009_repository_authorization.sql"),
+ include_str!("../src/store/migrations/010_audit_history.sql"),
+ "PRAGMA user_version = 10;\n",
+);
fn database(directory: &TempDir, name: &str) -> std::path::PathBuf {
directory.path().join(name)
@@ -150,7 +157,7 @@
let directory = TempDir::new().expect("create a temporary directory");
let store = Store::open(&database(&directory, "store.sqlite")).expect("open the store");
- assert_eq!(store.schema_version().expect("read the schema version"), 10);
+ assert_eq!(store.schema_version().expect("read the schema version"), 11);
assert_eq!(
store
.connection()
@@ -363,6 +370,109 @@
assert_eq!(imported_events[0].kind, "tag-created");
assert_eq!(imported_events[1].kind, "ref-created");
assert_eq!(imported_events[2].kind, "repository-imported");
+ assert_eq!(
+ imported_events
+ .iter()
+ .map(|event| event.sequence)
+ .collect::<Vec<_>>(),
+ vec![3, 2, 1]
+ );
+ for event in &imported_events {
+ assert_eq!(event.event_id.len(), 32);
+ assert!(
+ event
+ .event_id
+ .bytes()
+ .all(|byte| byte.is_ascii_digit() || (b'a'..=b'f').contains(&byte))
+ );
+ assert_eq!(event.payload_version, 1);
+ let payload: serde_json::Value =
+ serde_json::from_str(&event.payload).expect("parse a versioned event payload");
+ assert_eq!(payload["version"], 1);
+ }
+ let imported_payload: serde_json::Value = serde_json::from_str(&imported_events[2].payload)
+ .expect("parse the repository import payload");
+ assert_eq!(imported_payload["owner"], "alice");
+ assert_eq!(imported_payload["repository"], "project");
+ assert_eq!(imported_payload["object_format"], "sha256");
+ let event_plan: String = store
+ .connection()
+ .query_row(
+ "EXPLAIN QUERY PLAN
+ SELECT event_id FROM repository_event
+ WHERE repository_id = ?1 AND sequence < ?2
+ ORDER BY sequence DESC LIMIT ?3",
+ rusqlite::params![repository.id, 100, 20],
+ |row| row.get(3),
+ )
+ .expect("read the repository event query plan");
+ assert!(
+ event_plan.contains("repository_event_feed"),
+ "query plan: {event_plan}"
+ );
+ let duplicate_sequence = store
+ .connection()
+ .execute(
+ "INSERT INTO repository_event
+ (event_id, repository_id, sequence, kind, actor, payload_version,
+ payload, created_at)
+ VALUES ('aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa', ?1, 1, 'push', 'alice',
+ 1, '{\"version\":1}', 20)",
+ [repository.id],
+ )
+ .expect_err("reject a duplicate repository event sequence");
+ assert_eq!(
+ duplicate_sequence.sqlite_error_code(),
+ Some(ErrorCode::ConstraintViolation)
+ );
+ let unversioned_payload = store
+ .connection()
+ .execute(
+ "INSERT INTO repository_event
+ (event_id, repository_id, sequence, kind, actor, payload_version,
+ payload, created_at)
+ VALUES ('bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb', ?1, 4, 'push', 'alice',
+ 1, '{}', 20)",
+ [repository.id],
+ )
+ .expect_err("reject an unversioned repository event payload");
+ assert_eq!(
+ unversioned_payload.sqlite_error_code(),
+ Some(ErrorCode::ConstraintViolation)
+ );
+ store
+ .connection()
+ .execute_batch(
+ "CREATE TEMP TRIGGER reject_repository_event
+ BEFORE INSERT ON repository_event
+ BEGIN
+ SELECT RAISE(ABORT, 'injected event failure');
+ END;",
+ )
+ .expect("install an event failure trigger");
+ let rejected_repository = NewRepository {
+ id: "11111111111111111111111111111111",
+ owner: "alice",
+ slug: "event-failure",
+ object_format: "sha1",
+ created_at: 20,
+ origin: RepositoryOrigin::Created,
+ initial_references: &[],
+ actor: "alice",
+ correlation_id: "test-event-failure",
+ };
+ assert!(matches!(
+ store.create_repository(&rejected_repository),
+ Err(StoreError::Sqlite(_))
+ ));
+ assert!(matches!(
+ store.repository("alice", "event-failure"),
+ Err(StoreError::RepositoryNotFound(_, _))
+ ));
+ store
+ .connection()
+ .execute_batch("DROP TRIGGER reject_repository_event;")
+ .expect("remove the event failure trigger");
let initial = format!("{} refs/heads/main\n", "0".repeat(64));
let proposed = format!("{} refs/heads/main\n", "a".repeat(64));
@@ -389,6 +499,33 @@
.expect("read events while a push is promoted");
assert_eq!(promoted_events.len(), 3);
store
+ .connection()
+ .execute_batch(
+ "CREATE TEMP TRIGGER reject_push_event
+ BEFORE INSERT ON repository_event
+ BEGIN
+ SELECT RAISE(ABORT, 'injected push event failure');
+ END;",
+ )
+ .expect("install a push event failure trigger");
+ assert!(matches!(
+ store.complete_git_intent(push.id),
+ Err(StoreError::Sqlite(_))
+ ));
+ assert!(
+ !store
+ .git_intent_completed(push.id)
+ .expect("read the rolled-back intent state")
+ );
+ let (_, rolled_back_events) = store
+ .public_repository_events("alice", "project", None, 10)
+ .expect("read events after the failed completion");
+ assert_eq!(rolled_back_events.len(), 3);
+ store
+ .connection()
+ .execute_batch("DROP TRIGGER reject_push_event;")
+ .expect("remove the push event failure trigger");
+ store
.complete_git_intent(push.id)
.expect("complete a managed push");
assert!(
@@ -400,12 +537,49 @@
.public_repository_events("alice", "project", None, 10)
.expect("read pushed repository events");
assert_eq!(pushed_events.len(), 5);
+ assert_eq!(
+ pushed_events
+ .iter()
+ .map(|event| event.sequence)
+ .collect::<Vec<_>>(),
+ vec![5, 4, 3, 2, 1]
+ );
assert_eq!(pushed_events[0].kind, "ref-created");
assert_eq!(
pushed_events[0].new_target.as_deref(),
Some("aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa")
);
assert_eq!(pushed_events[1].kind, "push");
+ let push_payload: serde_json::Value =
+ serde_json::from_str(&pushed_events[1].payload).expect("parse the push payload");
+ assert_eq!(push_payload["operation_id"], push.id);
+ let (_, older_events) = store
+ .public_repository_events("alice", "project", Some(4), 10)
+ .expect("read events before a repository sequence");
+ assert_eq!(
+ older_events
+ .iter()
+ .map(|event| event.sequence)
+ .collect::<Vec<_>>(),
+ vec![3, 2, 1]
+ );
+ let event_ids = pushed_events
+ .iter()
+ .map(|event| event.event_id.clone())
+ .collect::<Vec<_>>();
+ let reopened = Store::open(&database(&directory, "store.sqlite"))
+ .expect("reopen the repository event store");
+ let (_, reopened_events) = reopened
+ .public_repository_events("alice", "project", None, 10)
+ .expect("read repository events after a reopen");
+ assert_eq!(
+ reopened_events
+ .iter()
+ .map(|event| event.event_id.clone())
+ .collect::<Vec<_>>(),
+ event_ids
+ );
+ drop(reopened);
let audits = store.audit_events(10).expect("read audit history");
assert_eq!(audits.len(), 2);
assert_eq!(audits[0].action, "ref.update");
@@ -459,6 +633,10 @@
store.create_repository(&missing_owner),
Err(StoreError::AccountNotFound(owner)) if owner == "bob"
));
+ let (_, unchanged_events) = store
+ .public_repository_events("alice", "project", None, 10)
+ .expect("read events after rejected repository mutations");
+ assert_eq!(unchanged_events.len(), 5);
store
.rename_repository(
@@ -836,13 +1014,14 @@
(V7_FIXTURE, 7),
(V8_FIXTURE, 8),
(V9_FIXTURE, 9),
+ (V10_FIXTURE, 10),
] {
let directory = TempDir::new().expect("create a temporary directory");
let path = database(&directory, "tit.sqlite3");
create_fixture(&path, fixture);
let store = Store::open(&path).expect("migrate the fixture");
- assert_eq!(store.schema_version().expect("read the schema version"), 10);
+ assert_eq!(store.schema_version().expect("read the schema version"), 11);
store.integrity_check().expect("check migrated integrity");
let state: String = store
.connection()
@@ -903,12 +1082,19 @@
.expect("read the backfilled event");
assert_eq!(events.len(), 1);
assert_eq!(events[0].kind, "repository-created");
+ assert_eq!(events[0].sequence, 1);
+ assert_eq!(events[0].payload_version, 1);
+ assert_eq!(
+ serde_json::from_str::<serde_json::Value>(&events[0].payload)
+ .expect("parse the backfilled event payload")["version"],
+ 1
+ );
assert_eq!(events[0].created_at, 2);
}
#[test]
fn recovers_complete_schema_versions_after_a_process_kill_during_migration() {
- for (mode, expected_version) in [("migration-uncommitted", 1), ("migration-committed", 10)] {
+ for (mode, expected_version) in [("migration-uncommitted", 1), ("migration-committed", 11)] {
let directory = TempDir::new().expect("create a temporary directory");
let path = database(&directory, "fixture.sqlite");
create_fixture(&path, V1_FIXTURE);