#![allow(
dead_code,
clippy::unwrap_used,
clippy::expect_used,
clippy::panic,
clippy::print_stdout
)]
use std::time::Duration;
use chrono::{DateTime, TimeDelta, Utc};
use turnframe_core::case::CaseRef;
use turnframe_core::command::{CommandOrigin, IdempotencyKey};
use turnframe_core::event::{CommittedEvent, OutboxEntry, OutboxStatus};
use turnframe_core::ids::{
AccountId, CaseRevision, CommandId, ConversationId, EventId, InteractionId, OutboxId, TurnId,
UserId,
};
use turnframe_core::interaction::{
Interaction, InteractionKind, InteractionOption, InteractionPayload, InteractionSpec,
StoredInteractionAction,
};
use turnframe_core::locale::Locale;
use turnframe_core::turn::{ActorContext, TurnInput};
use turnframe_store::conversation::StoredUserTurn;
use turnframe_store::events::EventBatch;
use turnframe_store::journal::{CommandJournalEntry, CommandJournalStatus};
use turnframe_store_postgres::{PgStoreConfig, PgStores};
use uuid::Uuid;
pub const DATABASE_URL_VARIABLE: &str = "TURNFRAME_TEST_DATABASE_URL";
const TABLES: &str = "tf_conversation, tf_turn, tf_turn_phase, tf_interaction, \
tf_command_journal, tf_domain_event, tf_outbox, tf_replay";
fn load_dotenv() {
static LOADED: std::sync::Once = std::sync::Once::new();
LOADED.call_once(|| {
let _ = dotenvy::from_path(concat!(env!("CARGO_MANIFEST_DIR"), "/../../.env"));
});
}
pub fn database_url() -> Option<String> {
load_dotenv();
std::env::var(DATABASE_URL_VARIABLE)
.ok()
.filter(|url| !url.trim().is_empty())
}
pub fn skipped(test: &str) {
println!("SKIPPED {test}: {DATABASE_URL_VARIABLE} is not set, so no database was available");
}
pub async fn sole_owner_of(url: &str, schema: &str) -> PgStores {
let store = migrated(url, schema).await;
truncate_all(&store).await;
store
}
pub async fn co_tenant_of(url: &str, schema: &str) -> PgStores {
migrated(url, schema).await
}
async fn migrated(url: &str, schema: &str) -> PgStores {
let store = open(url, schema).await;
store.migrate().await.expect("migrations apply");
store
}
pub async fn second_pool(url: &str, schema: &str) -> PgStores {
open(url, schema).await
}
async fn open(url: &str, schema: &str) -> PgStores {
let config = PgStoreConfig::new()
.max_connections(8)
.min_connections(0)
.acquire_timeout(Duration::from_secs(30))
.schema(schema)
.expect("the schema name is a plain identifier");
PgStores::connect_with(url, &config)
.await
.expect("the test database accepts a connection")
}
pub async fn truncate_all(store: &PgStores) {
let statement = format!("TRUNCATE {TABLES} RESTART IDENTITY CASCADE");
sqlx::query(&statement)
.execute(store.pool())
.await
.expect("the test schema can be emptied");
}
pub fn unique_account(test: &str) -> AccountId {
AccountId::from(format!("{test}-{}", Uuid::new_v4().simple()))
}
pub fn epoch() -> DateTime<Utc> {
DateTime::<Utc>::UNIX_EPOCH
}
pub fn at(seconds: i64) -> DateTime<Utc> {
epoch() + TimeDelta::seconds(seconds)
}
pub fn case(case_id: &str, revision: u64) -> CaseRef {
CaseRef::new("turnframe-postgres-test", case_id, CaseRevision(revision))
}
pub fn card(account: &AccountId, case_ref: CaseRef, created_at: DateTime<Utc>) -> Interaction {
let payload = InteractionPayload::new("a card").with_option(InteractionOption::new(
"ack",
"Got it",
StoredInteractionAction::Dismiss,
));
let spec = InteractionSpec::new("card", case_ref, InteractionKind::SingleSelect, payload);
Interaction::from_spec(
spec,
InteractionId::new(),
account.clone(),
ConversationId::new(),
TurnId::new(),
created_at,
)
.expect("the fixture card is answerable")
}
pub fn user_turn(
account: &AccountId,
conversation: ConversationId,
turn: TurnId,
received_at: DateTime<Utc>,
) -> StoredUserTurn {
StoredUserTurn::new(
TurnInput {
turn_id: turn,
conversation_id: conversation,
actor: ActorContext::new(account.clone(), UserId::from("test-user")),
text: Some("a turn".to_owned()),
interaction_response: None,
attachments: Vec::new(),
origin: None,
locale: Locale::from("en"),
effort: None,
},
received_at,
)
}
pub fn journal_entry(
account: &AccountId,
turn: TurnId,
command_id: CommandId,
key: &str,
) -> CommandJournalEntry {
CommandJournalEntry {
command_id,
account_id: account.clone(),
idempotency_key: IdempotencyKey::new(key),
turn_id: turn,
case_ref: case("case-1", 1),
command_type: "turnframe.postgres.test".to_owned(),
command_payload: serde_json::json!({ "key": key }),
origin: CommandOrigin::InternalPolicy {
policy_key: "test".to_owned(),
},
status: CommandJournalStatus::Pending,
result: None,
created_at: epoch(),
completed_at: None,
}
}
pub fn event_batch(
account: &AccountId,
case_id: &str,
command_id: CommandId,
revision: u64,
ids: &[EventId],
) -> EventBatch {
EventBatch::new(
account.clone(),
case(case_id, 0).key(),
command_id,
CaseRevision(revision),
ids.iter()
.map(|id| CommittedEvent {
event_id: *id,
event_type: "turnframe.postgres.test.happened".to_owned(),
occurred_at: epoch(),
payload: serde_json::json!({}),
})
.collect(),
)
}
pub fn outbox_entry(
outbox_id: OutboxId,
command_id: CommandId,
destination: &str,
key: &str,
created_at: DateTime<Utc>,
) -> OutboxEntry {
OutboxEntry {
outbox_id,
command_id,
destination: destination.to_owned(),
payload: serde_json::json!({ "key": key }),
idempotency_key: IdempotencyKey::new(key),
status: OutboxStatus::Pending,
attempt_count: 0,
next_attempt_at: None,
created_at,
completed_at: None,
}
}