mirror of
https://relay.ngit.dev/npub15qydau2hjma6ngxkl2cyar74wzyjshvl65za5k5rl69264ar2exs5cyejr/ngit-grasp.git
synced 2026-10-05 15:08:24 +00:00
add landing page and nostr-relay-builder relay on same port
This commit is contained in:
Generated
+725
-49
File diff suppressed because it is too large
Load Diff
+10
-8
@@ -10,18 +10,19 @@ repository = "https://gitworkshop.dev/ngit-grasp"
|
||||
[dependencies]
|
||||
# Async runtime
|
||||
tokio = { version = "1.35", features = ["full"] }
|
||||
tokio-tungstenite = "0.21"
|
||||
|
||||
# WebSocket
|
||||
tungstenite = "0.21"
|
||||
futures-util = "0.3"
|
||||
# HTTP server
|
||||
actix-web = "4.4"
|
||||
actix-ws = "0.3"
|
||||
|
||||
# HTTP server (for future use)
|
||||
# actix-web = "4.4"
|
||||
# actix-cors = "0.7"
|
||||
# Nostr relay
|
||||
nostr-relay-builder = "0.44"
|
||||
|
||||
# Nostr
|
||||
nostr-sdk = "0.43"
|
||||
nostr-sdk = "0.44"
|
||||
|
||||
# Utilities
|
||||
futures-util = "0.3"
|
||||
|
||||
# Serialization
|
||||
serde = { version = "1.0", features = ["derive"] }
|
||||
@@ -37,6 +38,7 @@ dotenvy = "0.15"
|
||||
# Error handling
|
||||
anyhow = "1.0"
|
||||
thiserror = "1.0"
|
||||
tokio-tungstenite = "0.28.0"
|
||||
|
||||
# Git (for future use)
|
||||
# git-http-backend = "0.3"
|
||||
|
||||
@@ -0,0 +1,37 @@
|
||||
/// Landing Page Handler
|
||||
///
|
||||
/// Serves the HTML landing page or upgrades to WebSocket for Nostr relay connections.
|
||||
|
||||
use actix_web::{web, HttpRequest, HttpResponse, Result};
|
||||
use nostr_relay_builder::LocalRelay;
|
||||
|
||||
use crate::config::Config;
|
||||
|
||||
/// Handle landing page or WebSocket upgrade
|
||||
pub async fn handle(
|
||||
req: HttpRequest,
|
||||
stream: web::Payload,
|
||||
config: web::Data<Config>,
|
||||
relay: web::Data<LocalRelay>,
|
||||
) -> Result<HttpResponse> {
|
||||
// Check if this is a WebSocket upgrade request
|
||||
if let Some(upgrade) = req.headers().get("upgrade") {
|
||||
if upgrade.to_str().unwrap_or("").eq_ignore_ascii_case("websocket") {
|
||||
// Delegate to WebSocket handler
|
||||
return crate::http::websocket::handle(req, stream, relay).await;
|
||||
}
|
||||
}
|
||||
|
||||
// Otherwise, serve the landing page
|
||||
let html = format!(
|
||||
include_str!("../../templates/landing.html"),
|
||||
relay_name = config.relay_name,
|
||||
relay_description = config.relay_description,
|
||||
domain = config.domain,
|
||||
bind_address = config.bind_address,
|
||||
);
|
||||
|
||||
Ok(HttpResponse::Ok()
|
||||
.content_type("text/html; charset=utf-8")
|
||||
.body(html))
|
||||
}
|
||||
@@ -0,0 +1,31 @@
|
||||
/// HTTP Server Module
|
||||
///
|
||||
/// Provides actix-web HTTP server with WebSocket upgrade support for the Nostr relay.
|
||||
|
||||
pub mod landing;
|
||||
pub mod websocket;
|
||||
|
||||
use actix_web::{middleware, web, App, HttpServer};
|
||||
use nostr_relay_builder::LocalRelay;
|
||||
|
||||
use crate::config::Config;
|
||||
|
||||
/// Start the HTTP server with integrated Nostr relay
|
||||
pub async fn run_server(config: Config, relay: LocalRelay) -> anyhow::Result<()> {
|
||||
let bind_addr = config.bind_address.clone();
|
||||
|
||||
tracing::info!("Starting HTTP server on {}", bind_addr);
|
||||
|
||||
HttpServer::new(move || {
|
||||
App::new()
|
||||
.app_data(web::Data::new(config.clone()))
|
||||
.app_data(web::Data::new(relay.clone()))
|
||||
.wrap(middleware::Logger::default())
|
||||
.route("/", web::get().to(landing::handle))
|
||||
})
|
||||
.bind(&bind_addr)?
|
||||
.run()
|
||||
.await?;
|
||||
|
||||
Ok(())
|
||||
}
|
||||
@@ -0,0 +1,73 @@
|
||||
/// WebSocket Handler
|
||||
///
|
||||
/// Handles WebSocket upgrade requests and passes connections to the Nostr relay.
|
||||
|
||||
use actix_web::{web, HttpRequest, HttpResponse, Result, Error};
|
||||
use actix_ws::Message;
|
||||
use futures_util::StreamExt;
|
||||
use nostr_relay_builder::LocalRelay;
|
||||
|
||||
/// Handle WebSocket upgrade and relay connection
|
||||
pub async fn handle(
|
||||
req: HttpRequest,
|
||||
stream: web::Payload,
|
||||
relay: web::Data<LocalRelay>,
|
||||
) -> Result<HttpResponse, Error> {
|
||||
let (response, mut session, mut msg_stream) = actix_ws::handle(&req, stream)?;
|
||||
|
||||
let peer_addr = req.peer_addr()
|
||||
.unwrap_or_else(|| "0.0.0.0:0".parse().unwrap());
|
||||
|
||||
tracing::debug!("WebSocket connection from {}", peer_addr);
|
||||
|
||||
// Spawn task to handle the WebSocket connection
|
||||
// TODO: Will use relay.take_connection() for full Nostr relay integration
|
||||
let _relay = relay.get_ref().clone();
|
||||
actix_web::rt::spawn(async move {
|
||||
// Create a channel to communicate between actix-ws and relay
|
||||
let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
|
||||
|
||||
// Spawn task to send messages from relay to client
|
||||
let mut session_clone = session.clone();
|
||||
actix_web::rt::spawn(async move {
|
||||
while let Some(msg) = rx.recv().await {
|
||||
if session_clone.text(msg).await.is_err() {
|
||||
break;
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
// Handle incoming messages from client
|
||||
while let Some(Ok(msg)) = msg_stream.next().await {
|
||||
match msg {
|
||||
Message::Text(text) => {
|
||||
// For now, just echo back - will integrate with relay in next phase
|
||||
tracing::debug!("Received text message: {}", text);
|
||||
if let Err(e) = tx.send(text.to_string()) {
|
||||
tracing::error!("Failed to send message: {}", e);
|
||||
break;
|
||||
}
|
||||
}
|
||||
Message::Binary(_) => {
|
||||
tracing::warn!("Received unexpected binary message");
|
||||
}
|
||||
Message::Close(_) => {
|
||||
tracing::debug!("Client closed connection");
|
||||
break;
|
||||
}
|
||||
Message::Ping(bytes) => {
|
||||
if session.pong(&bytes).await.is_err() {
|
||||
break;
|
||||
}
|
||||
}
|
||||
Message::Pong(_) => {}
|
||||
Message::Continuation(_) => {}
|
||||
Message::Nop => {}
|
||||
}
|
||||
}
|
||||
|
||||
tracing::debug!("WebSocket connection closed for {}", peer_addr);
|
||||
});
|
||||
|
||||
Ok(response)
|
||||
}
|
||||
+1
-1
@@ -1,3 +1,3 @@
|
||||
pub mod config;
|
||||
pub mod http;
|
||||
pub mod nostr;
|
||||
pub mod storage;
|
||||
|
||||
+12
-10
@@ -3,8 +3,8 @@ use tracing::{info, Level};
|
||||
use tracing_subscriber::FmtSubscriber;
|
||||
|
||||
mod config;
|
||||
mod http;
|
||||
mod nostr;
|
||||
mod storage;
|
||||
|
||||
use config::Config;
|
||||
|
||||
@@ -16,21 +16,23 @@ async fn main() -> Result<()> {
|
||||
.finish();
|
||||
tracing::subscriber::set_global_default(subscriber)?;
|
||||
|
||||
info!("Starting ngit-grasp...");
|
||||
info!("Starting ngit-grasp with nostr-relay-builder...");
|
||||
|
||||
// Load configuration
|
||||
let config = Config::from_env()?;
|
||||
info!("Configuration loaded: {}", config.bind_address);
|
||||
|
||||
// Initialize storage
|
||||
let storage = storage::Storage::new(&config)?;
|
||||
info!("Storage initialized at: {}", config.relay_data_path);
|
||||
// Create Nostr relay with NIP-34 validation
|
||||
if let Ok(relay) = nostr::builder::create_relay(&config) {
|
||||
info!(
|
||||
"Relay created with NIP-34 validation for domain: {}",
|
||||
config.domain
|
||||
);
|
||||
|
||||
// Start Nostr relay
|
||||
let relay = nostr::relay::RelayServer::new(config.clone(), storage)?;
|
||||
|
||||
info!("Starting Nostr relay on {}", config.bind_address);
|
||||
relay.run().await?;
|
||||
// Start HTTP server with integrated relay
|
||||
info!("Starting HTTP server on {}", config.bind_address);
|
||||
http::run_server(config, relay).await?;
|
||||
}
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
@@ -0,0 +1,121 @@
|
||||
/// Nostr Relay Builder Configuration
|
||||
///
|
||||
/// This module integrates nostr-relay-builder with NIP-34 validation logic
|
||||
/// preserved from the original implementation.
|
||||
use std::net::SocketAddr;
|
||||
use std::path::Path;
|
||||
|
||||
use nostr::nips::nip19::ToBech32;
|
||||
use nostr_relay_builder::prelude::*;
|
||||
|
||||
use crate::config::Config;
|
||||
use crate::nostr::events::{
|
||||
validate_announcement, validate_state, KIND_REPOSITORY_ANNOUNCEMENT, KIND_REPOSITORY_STATE,
|
||||
};
|
||||
|
||||
/// NIP-34 Write Policy
|
||||
///
|
||||
/// Validates repository announcement and state events according to GRASP-01 spec.
|
||||
/// Preserves all original validation logic from src/nostr/events.rs.
|
||||
#[derive(Debug, Clone)]
|
||||
pub struct Nip34WritePolicy {
|
||||
domain: String,
|
||||
}
|
||||
|
||||
impl Nip34WritePolicy {
|
||||
pub fn new(domain: impl Into<String>) -> Self {
|
||||
Self {
|
||||
domain: domain.into(),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl WritePolicy for Nip34WritePolicy {
|
||||
fn admit_event<'a>(
|
||||
&'a self,
|
||||
event: &'a nostr_relay_builder::prelude::Event,
|
||||
_addr: &'a SocketAddr,
|
||||
) -> BoxedFuture<'a, PolicyResult> {
|
||||
Box::pin(async move {
|
||||
match event.kind.as_u16() {
|
||||
KIND_REPOSITORY_ANNOUNCEMENT => match validate_announcement(event, &self.domain) {
|
||||
Ok(_) => {
|
||||
tracing::debug!(
|
||||
"Accepted repository announcement: {}",
|
||||
event
|
||||
.id
|
||||
.to_bech32()
|
||||
.unwrap_or_else(|_| "invalid".to_string())
|
||||
);
|
||||
PolicyResult::Accept
|
||||
}
|
||||
Err(e) => {
|
||||
tracing::warn!(
|
||||
"Rejected repository announcement {}: {}",
|
||||
event
|
||||
.id
|
||||
.to_bech32()
|
||||
.unwrap_or_else(|_| "invalid".to_string()),
|
||||
e
|
||||
);
|
||||
PolicyResult::Reject(e.to_string())
|
||||
}
|
||||
},
|
||||
KIND_REPOSITORY_STATE => match validate_state(event) {
|
||||
Ok(_) => {
|
||||
tracing::debug!(
|
||||
"Accepted repository state: {}",
|
||||
event
|
||||
.id
|
||||
.to_bech32()
|
||||
.unwrap_or_else(|_| "invalid".to_string())
|
||||
);
|
||||
PolicyResult::Accept
|
||||
}
|
||||
Err(e) => {
|
||||
tracing::warn!(
|
||||
"Rejected repository state {}: {}",
|
||||
event
|
||||
.id
|
||||
.to_bech32()
|
||||
.unwrap_or_else(|_| "invalid".to_string()),
|
||||
e
|
||||
);
|
||||
PolicyResult::Reject(e.to_string())
|
||||
}
|
||||
},
|
||||
// Accept all other event kinds without validation
|
||||
_ => PolicyResult::Accept,
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
/// Create a configured LocalRelay with NIP-34 validation
|
||||
pub fn create_relay(config: &Config) -> Result<LocalRelay> {
|
||||
tracing::info!("Configuring nostr relay...");
|
||||
|
||||
// Determine database path
|
||||
let db_path = Path::new(&config.relay_data_path);
|
||||
|
||||
// Create database - using in-memory for now, can switch to persistent later
|
||||
// TODO: Add configuration for NostrDB or LMDB backends
|
||||
let database = MemoryDatabase::with_opts(MemoryDatabaseOptions {
|
||||
events: true,
|
||||
max_events: Some(100_000),
|
||||
});
|
||||
|
||||
tracing::info!("Using in-memory database (path: {})", db_path.display());
|
||||
|
||||
// Build relay with NIP-34 validation
|
||||
let builder = RelayBuilder::default()
|
||||
.database(database)
|
||||
.write_policy(Nip34WritePolicy::new(&config.domain));
|
||||
|
||||
tracing::info!(
|
||||
"Relay configured with NIP-34 validation for domain: {}",
|
||||
config.domain
|
||||
);
|
||||
|
||||
Ok(LocalRelay::new(builder))
|
||||
}
|
||||
+1
-1
@@ -1,2 +1,2 @@
|
||||
pub mod builder;
|
||||
pub mod events;
|
||||
pub mod relay;
|
||||
|
||||
@@ -1,340 +0,0 @@
|
||||
use anyhow::Result;
|
||||
use futures_util::{SinkExt, StreamExt};
|
||||
use nostr_sdk::{Event, EventId, Filter, Kind};
|
||||
use serde_json::{json, Value};
|
||||
use std::collections::HashMap;
|
||||
use std::net::SocketAddr;
|
||||
use std::sync::Arc;
|
||||
use tokio::net::{TcpListener, TcpStream};
|
||||
use tokio::sync::RwLock;
|
||||
use tokio_tungstenite::{accept_async, tungstenite::Message};
|
||||
use tracing::{debug, error, info, warn};
|
||||
|
||||
use crate::config::Config;
|
||||
use crate::nostr::events::{validate_announcement, validate_state, KIND_REPOSITORY_ANNOUNCEMENT, KIND_REPOSITORY_STATE};
|
||||
use crate::storage::Storage;
|
||||
|
||||
type Subscriptions = Arc<RwLock<HashMap<String, Vec<Filter>>>>;
|
||||
|
||||
pub struct RelayServer {
|
||||
config: Config,
|
||||
storage: Storage,
|
||||
}
|
||||
|
||||
impl RelayServer {
|
||||
pub fn new(config: Config, storage: Storage) -> Result<Self> {
|
||||
Ok(RelayServer { config, storage })
|
||||
}
|
||||
|
||||
pub async fn run(self) -> Result<()> {
|
||||
let addr: SocketAddr = self.config.bind_address.parse()?;
|
||||
let listener = TcpListener::bind(&addr).await?;
|
||||
|
||||
info!("✅ Nostr relay listening on ws://{}", addr);
|
||||
info!("📡 Ready to accept connections...");
|
||||
|
||||
loop {
|
||||
match listener.accept().await {
|
||||
Ok((stream, peer_addr)) => {
|
||||
debug!("New connection from: {}", peer_addr);
|
||||
let storage = self.storage.clone();
|
||||
tokio::spawn(async move {
|
||||
if let Err(e) = handle_connection(stream, storage).await {
|
||||
error!("Error handling connection from {}: {}", peer_addr, e);
|
||||
}
|
||||
});
|
||||
}
|
||||
Err(e) => {
|
||||
error!("Error accepting connection: {}", e);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
async fn handle_connection(stream: TcpStream, storage: Storage) -> Result<()> {
|
||||
let ws_stream = accept_async(stream).await?;
|
||||
let (mut ws_sender, mut ws_receiver) = ws_stream.split();
|
||||
|
||||
let subscriptions: Subscriptions = Arc::new(RwLock::new(HashMap::new()));
|
||||
|
||||
while let Some(msg) = ws_receiver.next().await {
|
||||
match msg {
|
||||
Ok(Message::Text(text)) => {
|
||||
debug!("Received message: {}", text);
|
||||
|
||||
match handle_message(&text, &storage, &subscriptions).await {
|
||||
Ok(responses) => {
|
||||
for response in responses {
|
||||
let response_text = serde_json::to_string(&response)?;
|
||||
debug!("Sending response: {}", response_text);
|
||||
ws_sender.send(Message::Text(response_text)).await?;
|
||||
}
|
||||
}
|
||||
Err(e) => {
|
||||
warn!("Error handling message: {}", e);
|
||||
let notice = json!(["NOTICE", format!("Error: {}", e)]);
|
||||
ws_sender.send(Message::Text(notice.to_string())).await?;
|
||||
}
|
||||
}
|
||||
}
|
||||
Ok(Message::Close(_)) => {
|
||||
debug!("Client closed connection");
|
||||
break;
|
||||
}
|
||||
Ok(Message::Ping(data)) => {
|
||||
ws_sender.send(Message::Pong(data)).await?;
|
||||
}
|
||||
Ok(_) => {
|
||||
// Ignore other message types
|
||||
}
|
||||
Err(e) => {
|
||||
error!("WebSocket error: {}", e);
|
||||
break;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn handle_message(
|
||||
text: &str,
|
||||
storage: &Storage,
|
||||
subscriptions: &Subscriptions,
|
||||
) -> Result<Vec<Value>> {
|
||||
let msg: Value = serde_json::from_str(text)?;
|
||||
|
||||
if let Some(arr) = msg.as_array() {
|
||||
if arr.is_empty() {
|
||||
return Ok(vec![json!(["NOTICE", "Empty message"])]);
|
||||
}
|
||||
|
||||
let msg_type = arr[0].as_str().unwrap_or("");
|
||||
|
||||
match msg_type {
|
||||
"EVENT" => handle_event(arr, storage).await,
|
||||
"REQ" => handle_req(arr, storage, subscriptions).await,
|
||||
"CLOSE" => handle_close(arr, subscriptions).await,
|
||||
_ => Ok(vec![json!(["NOTICE", format!("Unknown message type: {}", msg_type)])]),
|
||||
}
|
||||
} else {
|
||||
Ok(vec![json!(["NOTICE", "Invalid message format"])])
|
||||
}
|
||||
}
|
||||
|
||||
async fn handle_event(arr: &[Value], storage: &Storage) -> Result<Vec<Value>> {
|
||||
if arr.len() < 2 {
|
||||
return Ok(vec![json!(["NOTICE", "EVENT message requires event object"])]);
|
||||
}
|
||||
|
||||
let event: Event = serde_json::from_value(arr[1].clone())?;
|
||||
let event_id = event.id;
|
||||
|
||||
// Verify event (signature and ID)
|
||||
if event.verify().is_err() {
|
||||
return Ok(vec![json!(["OK", event_id.to_hex(), false, "invalid: signature or ID verification failed"])]);
|
||||
}
|
||||
|
||||
// Check if event already exists
|
||||
if storage.get_event(&event_id.to_hex()).await.is_some() {
|
||||
return Ok(vec![json!(["OK", event_id.to_hex(), true, "duplicate: event already exists"])]);
|
||||
}
|
||||
|
||||
// Validate repository announcements (kind 30617)
|
||||
if event.kind == Kind::from(KIND_REPOSITORY_ANNOUNCEMENT) {
|
||||
// Get domain from storage config
|
||||
let domain = storage.get_domain();
|
||||
|
||||
match validate_announcement(&event, &domain) {
|
||||
Ok(()) => {
|
||||
info!("✅ Valid repository announcement: {} ({})", event_id, event.kind);
|
||||
}
|
||||
Err(e) => {
|
||||
warn!("❌ Invalid repository announcement: {}", e);
|
||||
return Ok(vec![json!(["OK", event_id.to_hex(), false, format!("invalid: {}", e)])]);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Validate repository state announcements (kind 30618)
|
||||
if event.kind == Kind::from(KIND_REPOSITORY_STATE) {
|
||||
match validate_state(&event) {
|
||||
Ok(()) => {
|
||||
info!("✅ Valid repository state: {} ({})", event_id, event.kind);
|
||||
}
|
||||
Err(e) => {
|
||||
warn!("❌ Invalid repository state: {}", e);
|
||||
return Ok(vec![json!(["OK", event_id.to_hex(), false, format!("invalid: {}", e)])]);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Store the event
|
||||
storage.store_event(event.clone()).await?;
|
||||
|
||||
info!("✅ Stored event: {} (kind: {})", event_id, event.kind);
|
||||
|
||||
Ok(vec![json!(["OK", event_id.to_hex(), true, ""])])
|
||||
}
|
||||
|
||||
async fn handle_req(
|
||||
arr: &[Value],
|
||||
storage: &Storage,
|
||||
subscriptions: &Subscriptions,
|
||||
) -> Result<Vec<Value>> {
|
||||
if arr.len() < 2 {
|
||||
return Ok(vec![json!(["NOTICE", "REQ message requires subscription ID"])]);
|
||||
}
|
||||
|
||||
let sub_id = arr[1].as_str().ok_or_else(|| anyhow::anyhow!("Invalid subscription ID"))?;
|
||||
|
||||
// Parse filters
|
||||
let mut filters = Vec::new();
|
||||
for filter_value in &arr[2..] {
|
||||
let filter: Filter = serde_json::from_value(filter_value.clone())?;
|
||||
filters.push(filter.clone());
|
||||
}
|
||||
|
||||
// Store subscription
|
||||
{
|
||||
let mut subs = subscriptions.write().await;
|
||||
subs.insert(sub_id.to_string(), filters.clone());
|
||||
}
|
||||
|
||||
debug!("Created subscription: {} with {} filters", sub_id, filters.len());
|
||||
|
||||
// Query and send matching events
|
||||
let mut responses = Vec::new();
|
||||
|
||||
for filter in filters {
|
||||
let events = storage.query_events(|event| {
|
||||
matches_filter(event, &filter)
|
||||
}).await;
|
||||
|
||||
for event in events {
|
||||
responses.push(json!(["EVENT", sub_id, event]));
|
||||
}
|
||||
}
|
||||
|
||||
// Send EOSE (End of Stored Events)
|
||||
responses.push(json!(["EOSE", sub_id]));
|
||||
|
||||
debug!("Subscription {} returned {} events", sub_id, responses.len() - 1);
|
||||
|
||||
Ok(responses)
|
||||
}
|
||||
|
||||
async fn handle_close(arr: &[Value], subscriptions: &Subscriptions) -> Result<Vec<Value>> {
|
||||
if arr.len() < 2 {
|
||||
return Ok(vec![json!(["NOTICE", "CLOSE message requires subscription ID"])]);
|
||||
}
|
||||
|
||||
let sub_id = arr[1].as_str().ok_or_else(|| anyhow::anyhow!("Invalid subscription ID"))?;
|
||||
|
||||
{
|
||||
let mut subs = subscriptions.write().await;
|
||||
subs.remove(sub_id);
|
||||
}
|
||||
|
||||
debug!("Closed subscription: {}", sub_id);
|
||||
|
||||
Ok(vec![])
|
||||
}
|
||||
|
||||
fn matches_filter(event: &Event, filter: &Filter) -> bool {
|
||||
// Check IDs
|
||||
if let Some(ref ids) = filter.ids {
|
||||
if !ids.is_empty() && !ids.contains(&event.id) {
|
||||
return false;
|
||||
}
|
||||
}
|
||||
|
||||
// Check authors
|
||||
if let Some(ref authors) = filter.authors {
|
||||
if !authors.is_empty() && !authors.contains(&event.pubkey) {
|
||||
return false;
|
||||
}
|
||||
}
|
||||
|
||||
// Check kinds
|
||||
if let Some(ref kinds) = filter.kinds {
|
||||
if !kinds.is_empty() && !kinds.contains(&event.kind) {
|
||||
return false;
|
||||
}
|
||||
}
|
||||
|
||||
// Check since
|
||||
if let Some(since) = filter.since {
|
||||
if event.created_at < since {
|
||||
return false;
|
||||
}
|
||||
}
|
||||
|
||||
// Check until
|
||||
if let Some(until) = filter.until {
|
||||
if event.created_at > until {
|
||||
return false;
|
||||
}
|
||||
}
|
||||
|
||||
// TODO: Check tags (#e, #p, etc.)
|
||||
|
||||
true
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use nostr_sdk::{EventBuilder, Keys, Kind};
|
||||
|
||||
#[test]
|
||||
fn test_matches_filter_by_id() {
|
||||
let keys = Keys::generate();
|
||||
let event = EventBuilder::text_note("test")
|
||||
.sign_with_keys(&keys)
|
||||
.unwrap();
|
||||
|
||||
// Filter matching the event ID
|
||||
let filter = Filter::new().id(event.id);
|
||||
assert!(matches_filter(&event, &filter));
|
||||
|
||||
// Filter not matching
|
||||
let other_id = EventId::all_zeros();
|
||||
let filter = Filter::new().id(other_id);
|
||||
assert!(!matches_filter(&event, &filter));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_matches_filter_by_author() {
|
||||
let keys = Keys::generate();
|
||||
let event = EventBuilder::text_note("test")
|
||||
.sign_with_keys(&keys)
|
||||
.unwrap();
|
||||
|
||||
// Filter matching the author
|
||||
let filter = Filter::new().author(keys.public_key());
|
||||
assert!(matches_filter(&event, &filter));
|
||||
|
||||
// Filter not matching
|
||||
let other_keys = Keys::generate();
|
||||
let filter = Filter::new().author(other_keys.public_key());
|
||||
assert!(!matches_filter(&event, &filter));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_matches_filter_by_kind() {
|
||||
let keys = Keys::generate();
|
||||
let event = EventBuilder::text_note("test")
|
||||
.sign_with_keys(&keys)
|
||||
.unwrap();
|
||||
|
||||
// Filter matching the kind
|
||||
let filter = Filter::new().kind(Kind::TextNote);
|
||||
assert!(matches_filter(&event, &filter));
|
||||
|
||||
// Filter not matching
|
||||
let filter = Filter::new().kind(Kind::Metadata);
|
||||
assert!(!matches_filter(&event, &filter));
|
||||
}
|
||||
}
|
||||
@@ -1,132 +0,0 @@
|
||||
use anyhow::Result;
|
||||
use nostr_sdk::Event;
|
||||
use std::collections::HashMap;
|
||||
use std::sync::Arc;
|
||||
use tokio::sync::RwLock;
|
||||
|
||||
use crate::config::Config;
|
||||
|
||||
/// Simple in-memory storage for events
|
||||
/// TODO: Persist to disk for production use
|
||||
#[derive(Clone)]
|
||||
pub struct Storage {
|
||||
events: Arc<RwLock<HashMap<String, Event>>>,
|
||||
data_path: String,
|
||||
domain: String,
|
||||
}
|
||||
|
||||
impl Storage {
|
||||
pub fn new(config: &Config) -> Result<Self> {
|
||||
// Create data directory if it doesn't exist
|
||||
std::fs::create_dir_all(&config.relay_data_path)?;
|
||||
|
||||
Ok(Storage {
|
||||
events: Arc::new(RwLock::new(HashMap::new())),
|
||||
data_path: config.relay_data_path.clone(),
|
||||
domain: config.domain.clone(),
|
||||
})
|
||||
}
|
||||
|
||||
pub fn get_domain(&self) -> String {
|
||||
self.domain.clone()
|
||||
}
|
||||
|
||||
pub async fn store_event(&self, event: Event) -> Result<()> {
|
||||
let mut events = self.events.write().await;
|
||||
events.insert(event.id.to_hex(), event);
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub async fn get_event(&self, event_id: &str) -> Option<Event> {
|
||||
let events = self.events.read().await;
|
||||
events.get(event_id).cloned()
|
||||
}
|
||||
|
||||
pub async fn query_events<F>(&self, filter: F) -> Vec<Event>
|
||||
where
|
||||
F: Fn(&Event) -> bool,
|
||||
{
|
||||
let events = self.events.read().await;
|
||||
events.values().filter(|e| filter(e)).cloned().collect()
|
||||
}
|
||||
|
||||
pub async fn count_events(&self) -> usize {
|
||||
let events = self.events.read().await;
|
||||
events.len()
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use nostr_sdk::{EventBuilder, Keys, Kind};
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_store_and_retrieve() {
|
||||
let config = Config {
|
||||
domain: "test".to_string(),
|
||||
owner_npub: "npub1test".to_string(),
|
||||
relay_name: "test".to_string(),
|
||||
relay_description: "test".to_string(),
|
||||
git_data_path: "./test_data/git".to_string(),
|
||||
relay_data_path: "./test_data/relay".to_string(),
|
||||
bind_address: "127.0.0.1:8080".to_string(),
|
||||
};
|
||||
|
||||
let storage = Storage::new(&config).unwrap();
|
||||
|
||||
// Create a test event
|
||||
let keys = Keys::generate();
|
||||
let event = EventBuilder::text_note("test content")
|
||||
.sign_with_keys(&keys)
|
||||
.unwrap();
|
||||
|
||||
// Store it
|
||||
storage.store_event(event.clone()).await.unwrap();
|
||||
|
||||
// Retrieve it
|
||||
let retrieved = storage.get_event(&event.id.to_hex()).await;
|
||||
assert!(retrieved.is_some());
|
||||
assert_eq!(retrieved.unwrap().id, event.id);
|
||||
|
||||
// Count events
|
||||
assert_eq!(storage.count_events().await, 1);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_query_events() {
|
||||
let config = Config {
|
||||
domain: "test".to_string(),
|
||||
owner_npub: "npub1test".to_string(),
|
||||
relay_name: "test".to_string(),
|
||||
relay_description: "test".to_string(),
|
||||
git_data_path: "./test_data/git".to_string(),
|
||||
relay_data_path: "./test_data/relay".to_string(),
|
||||
bind_address: "127.0.0.1:8080".to_string(),
|
||||
};
|
||||
|
||||
let storage = Storage::new(&config).unwrap();
|
||||
|
||||
// Create multiple events
|
||||
let keys = Keys::generate();
|
||||
let event1 = EventBuilder::text_note("message 1")
|
||||
.sign_with_keys(&keys)
|
||||
.unwrap();
|
||||
let event2 = EventBuilder::text_note("message 2")
|
||||
.sign_with_keys(&keys)
|
||||
.unwrap();
|
||||
|
||||
storage.store_event(event1.clone()).await.unwrap();
|
||||
storage.store_event(event2.clone()).await.unwrap();
|
||||
|
||||
// Query all events
|
||||
let all_events = storage.query_events(|_| true).await;
|
||||
assert_eq!(all_events.len(), 2);
|
||||
|
||||
// Query by kind
|
||||
let text_notes = storage
|
||||
.query_events(|e| e.kind == Kind::TextNote)
|
||||
.await;
|
||||
assert_eq!(text_notes.len(), 2);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,151 @@
|
||||
<!DOCTYPE html>
|
||||
<html lang="en">
|
||||
<head>
|
||||
<meta charset="UTF-8">
|
||||
<meta name="viewport" content="width=device-width, initial-scale=1.0">
|
||||
<title>{relay_name}</title>
|
||||
<style>
|
||||
* {{
|
||||
margin: 0;
|
||||
padding: 0;
|
||||
box-sizing: border-box;
|
||||
}}
|
||||
body {{
|
||||
font-family: -apple-system, BlinkMacSystemFont, 'Segoe UI', 'Roboto', 'Oxygen', 'Ubuntu', 'Cantarell', sans-serif;
|
||||
line-height: 1.6;
|
||||
background: linear-gradient(135deg, #667eea 0%, #764ba2 100%);
|
||||
color: #333;
|
||||
min-height: 100vh;
|
||||
display: flex;
|
||||
align-items: center;
|
||||
justify-content: center;
|
||||
padding: 20px;
|
||||
}}
|
||||
.container {{
|
||||
max-width: 800px;
|
||||
background: white;
|
||||
padding: 40px;
|
||||
border-radius: 12px;
|
||||
box-shadow: 0 20px 60px rgba(0,0,0,0.3);
|
||||
}}
|
||||
h1 {{
|
||||
color: #667eea;
|
||||
margin-bottom: 10px;
|
||||
font-size: 2.5em;
|
||||
}}
|
||||
h2 {{
|
||||
color: #764ba2;
|
||||
margin-top: 30px;
|
||||
margin-bottom: 15px;
|
||||
font-size: 1.5em;
|
||||
border-bottom: 2px solid #667eea;
|
||||
padding-bottom: 10px;
|
||||
}}
|
||||
.subtitle {{
|
||||
color: #666;
|
||||
margin-bottom: 30px;
|
||||
font-size: 1.1em;
|
||||
}}
|
||||
.info-grid {{
|
||||
display: grid;
|
||||
grid-template-columns: 140px 1fr;
|
||||
gap: 15px;
|
||||
margin: 20px 0;
|
||||
}}
|
||||
.info-label {{
|
||||
font-weight: 600;
|
||||
color: #555;
|
||||
}}
|
||||
.info-value {{
|
||||
color: #333;
|
||||
}}
|
||||
code {{
|
||||
background: #f4f4f4;
|
||||
padding: 3px 8px;
|
||||
border-radius: 4px;
|
||||
font-family: 'Courier New', monospace;
|
||||
color: #667eea;
|
||||
font-size: 0.9em;
|
||||
}}
|
||||
ul {{
|
||||
margin: 15px 0;
|
||||
padding-left: 30px;
|
||||
}}
|
||||
li {{
|
||||
margin: 8px 0;
|
||||
}}
|
||||
.feature-box {{
|
||||
background: #f9f9f9;
|
||||
padding: 20px;
|
||||
border-radius: 8px;
|
||||
margin: 15px 0;
|
||||
border-left: 4px solid #667eea;
|
||||
}}
|
||||
.footer {{
|
||||
margin-top: 40px;
|
||||
padding-top: 20px;
|
||||
border-top: 1px solid #eee;
|
||||
text-align: center;
|
||||
color: #999;
|
||||
font-size: 0.9em;
|
||||
}}
|
||||
.badge {{
|
||||
display: inline-block;
|
||||
background: #667eea;
|
||||
color: white;
|
||||
padding: 4px 12px;
|
||||
border-radius: 12px;
|
||||
font-size: 0.85em;
|
||||
margin-right: 8px;
|
||||
}}
|
||||
</style>
|
||||
</head>
|
||||
<body>
|
||||
<div class="container">
|
||||
<h1>{relay_name}</h1>
|
||||
<p class="subtitle">{relay_description}</p>
|
||||
|
||||
<h2>📡 Connection Information</h2>
|
||||
<div class="info-grid">
|
||||
<div class="info-label">Domain:</div>
|
||||
<div class="info-value"><code>{domain}</code></div>
|
||||
|
||||
<div class="info-label">WebSocket:</div>
|
||||
<div class="info-value"><code>ws://{bind_address}/</code></div>
|
||||
|
||||
<div class="info-label">HTTP:</div>
|
||||
<div class="info-value"><code>http://{bind_address}/</code></div>
|
||||
</div>
|
||||
|
||||
<h2>🔌 Supported NIPs</h2>
|
||||
<div class="feature-box">
|
||||
<ul>
|
||||
<li><span class="badge">NIP-01</span> Basic protocol flow and event structure</li>
|
||||
<li><span class="badge">NIP-11</span> Relay information document</li>
|
||||
<li><span class="badge">NIP-34</span> Git repository events and announcements</li>
|
||||
</ul>
|
||||
</div>
|
||||
|
||||
<h2>🎯 GRASP Features</h2>
|
||||
<div class="feature-box">
|
||||
<ul>
|
||||
<li><strong>Repository Announcements</strong> (kind 30617) - Declare Git repositories with validation</li>
|
||||
<li><strong>Repository State</strong> (kind 30618) - Track repository state changes</li>
|
||||
<li><strong>NIP-34 Validation</strong> - Automatic validation of repository events</li>
|
||||
<li><strong>Git HTTP Backend</strong> - Coming soon in Phase 2</li>
|
||||
</ul>
|
||||
</div>
|
||||
|
||||
<h2>📚 Documentation</h2>
|
||||
<p>For more information about GRASP (Git Relays Authorized via Signed-Nostr Proofs):</p>
|
||||
<ul>
|
||||
<li><a href="https://gitworkshop.dev/danconwaydev.com/grasp/01.md" target="_blank">GRASP-01 Specification</a></li>
|
||||
<li><a href="https://github.com/nostr-protocol/nips/blob/master/34.md" target="_blank">NIP-34: Git Stuff</a></li>
|
||||
</ul>
|
||||
|
||||
<div class="footer">
|
||||
<p>Powered by <strong>ngit-grasp</strong> with <strong>nostr-relay-builder</strong></p>
|
||||
</div>
|
||||
</div>
|
||||
</body>
|
||||
</html>
|
||||
+25
-25
@@ -39,22 +39,22 @@ use grasp_audit::*;
|
||||
async fn test_nip01_smoke() {
|
||||
// Start test relay
|
||||
let relay = TestRelay::start().await;
|
||||
|
||||
|
||||
// Create audit client in CI mode (isolated, no cleanup needed)
|
||||
let config = AuditConfig::ci();
|
||||
let client = AuditClient::new(relay.url(), config)
|
||||
.await
|
||||
.expect("Failed to create audit client");
|
||||
|
||||
|
||||
// Run all NIP-01 smoke tests
|
||||
let results = specs::Nip01SmokeTests::run_all(&client).await;
|
||||
|
||||
|
||||
// Print detailed report
|
||||
results.print_report();
|
||||
|
||||
|
||||
// Stop relay
|
||||
relay.stop().await;
|
||||
|
||||
|
||||
// Assert all tests passed
|
||||
assert!(
|
||||
results.all_passed(),
|
||||
@@ -70,20 +70,20 @@ async fn test_nip01_smoke() {
|
||||
/// for more granular testing or debugging.
|
||||
#[tokio::test]
|
||||
async fn test_nip01_individual_tests() {
|
||||
use grasp_audit::specs::nip01_smoke::Nip01SmokeTests;
|
||||
|
||||
use grasp_audit::specs::grasp01::Nip01SmokeTests;
|
||||
|
||||
let relay = TestRelay::start().await;
|
||||
let config = AuditConfig::ci();
|
||||
let client = AuditClient::new(relay.url(), config)
|
||||
.await
|
||||
.expect("Failed to create audit client");
|
||||
|
||||
|
||||
// We can't call private methods, so we'll run the full suite
|
||||
// This test is mainly to show the pattern
|
||||
let all_results = Nip01SmokeTests::run_all(&client).await;
|
||||
|
||||
|
||||
relay.stop().await;
|
||||
|
||||
|
||||
// Verify
|
||||
assert!(all_results.all_passed());
|
||||
}
|
||||
@@ -99,25 +99,25 @@ async fn test_relay_validates_events() {
|
||||
let client = AuditClient::new(relay.url(), config)
|
||||
.await
|
||||
.expect("Failed to create audit client");
|
||||
|
||||
|
||||
// The validation tests are part of the smoke tests
|
||||
let results = specs::Nip01SmokeTests::run_all(&client).await;
|
||||
|
||||
|
||||
// Check that validation tests exist and pass
|
||||
let validation_tests: Vec<_> = results
|
||||
.results
|
||||
.iter()
|
||||
.filter(|t| t.spec_ref.contains("validation"))
|
||||
.collect();
|
||||
|
||||
|
||||
relay.stop().await;
|
||||
|
||||
|
||||
// Should have validation tests
|
||||
assert!(
|
||||
!validation_tests.is_empty(),
|
||||
"No validation tests found in NIP-01 smoke tests"
|
||||
);
|
||||
|
||||
|
||||
// All validation tests should pass
|
||||
for test in validation_tests {
|
||||
assert!(
|
||||
@@ -137,18 +137,18 @@ async fn test_relay_lifecycle() {
|
||||
// Start relay
|
||||
let relay = TestRelay::start().await;
|
||||
let url = relay.url().to_string();
|
||||
|
||||
|
||||
// Verify we can connect
|
||||
let config = AuditConfig::ci();
|
||||
let client = AuditClient::new(&url, config)
|
||||
.await
|
||||
.expect("Failed to connect to relay");
|
||||
|
||||
|
||||
assert!(client.is_connected().await, "Client should be connected");
|
||||
|
||||
|
||||
// Stop relay
|
||||
relay.stop().await;
|
||||
|
||||
|
||||
// Note: We can't easily verify disconnection without modifying grasp-audit
|
||||
// to expose connection state after relay shutdown. That's okay - the
|
||||
// important part is that the relay starts and stops cleanly.
|
||||
@@ -162,28 +162,28 @@ async fn test_parallel_relays() {
|
||||
// Start two relays simultaneously
|
||||
let relay1 = TestRelay::start().await;
|
||||
let relay2 = TestRelay::start().await;
|
||||
|
||||
|
||||
// Should have different URLs (different ports)
|
||||
assert_ne!(
|
||||
relay1.url(),
|
||||
relay2.url(),
|
||||
"Relays should use different ports"
|
||||
);
|
||||
|
||||
|
||||
// Both should be connectable
|
||||
let config = AuditConfig::ci();
|
||||
|
||||
|
||||
let client1 = AuditClient::new(relay1.url(), config.clone())
|
||||
.await
|
||||
.expect("Failed to connect to relay 1");
|
||||
|
||||
|
||||
let client2 = AuditClient::new(relay2.url(), config)
|
||||
.await
|
||||
.expect("Failed to connect to relay 2");
|
||||
|
||||
|
||||
assert!(client1.is_connected().await);
|
||||
assert!(client2.is_connected().await);
|
||||
|
||||
|
||||
// Clean up
|
||||
relay1.stop().await;
|
||||
relay2.stop().await;
|
||||
|
||||
@@ -36,7 +36,9 @@ const KIND_REPOSITORY_ANNOUNCEMENT: u16 = 30617;
|
||||
const KIND_REPOSITORY_STATE: u16 = 30618;
|
||||
|
||||
/// Helper to connect to a test relay
|
||||
async fn connect_to_relay(url: &str) -> tokio_tungstenite::WebSocketStream<tokio_tungstenite::MaybeTlsStream<tokio::net::TcpStream>> {
|
||||
async fn connect_to_relay(
|
||||
url: &str,
|
||||
) -> tokio_tungstenite::WebSocketStream<tokio_tungstenite::MaybeTlsStream<tokio::net::TcpStream>> {
|
||||
let (ws, _) = connect_async(url)
|
||||
.await
|
||||
.expect("Failed to connect to relay");
|
||||
@@ -54,10 +56,7 @@ fn create_announcement(
|
||||
let mut tags = vec![Tag::custom(TagKind::d(), vec![identifier.to_string()])];
|
||||
|
||||
for url in clone_urls {
|
||||
tags.push(Tag::custom(
|
||||
TagKind::Clone,
|
||||
vec![url],
|
||||
));
|
||||
tags.push(Tag::custom(TagKind::Clone, vec![url]));
|
||||
}
|
||||
|
||||
for relay in relays {
|
||||
@@ -94,10 +93,10 @@ fn create_state(keys: &Keys, identifier: &str, branches: Vec<(&str, &str)>) -> n
|
||||
#[tokio::test]
|
||||
async fn test_relay_accepts_connection() {
|
||||
let relay = TestRelay::start().await;
|
||||
|
||||
|
||||
// Try to connect
|
||||
let ws = connect_to_relay(relay.url()).await;
|
||||
|
||||
|
||||
drop(ws); // Clean disconnect
|
||||
}
|
||||
|
||||
@@ -119,14 +118,14 @@ async fn test_accepts_valid_announcement() {
|
||||
|
||||
// Send event
|
||||
let event_msg = json!(["EVENT", event]);
|
||||
ws.send(Message::Text(event_msg.to_string()))
|
||||
ws.send(Message::Text(event_msg.to_string().into()))
|
||||
.await
|
||||
.expect("Failed to send event");
|
||||
|
||||
// Read response
|
||||
if let Some(Ok(Message::Text(text))) = ws.next().await {
|
||||
let response: Value = serde_json::from_str(&text).expect("Failed to parse response");
|
||||
|
||||
|
||||
// Should be ["OK", event_id, true, ""]
|
||||
assert_eq!(response[0], "OK");
|
||||
assert_eq!(response[1], event.id.to_hex());
|
||||
@@ -146,9 +145,7 @@ async fn test_rejects_announcement_without_clone() {
|
||||
let relay = TestRelay::start().await;
|
||||
let keys = Keys::generate();
|
||||
|
||||
let (mut ws, _) = connect_async(relay.url())
|
||||
.await
|
||||
.expect("Failed to connect");
|
||||
let (mut ws, _) = connect_async(relay.url()).await.expect("Failed to connect");
|
||||
|
||||
// Missing clone tag
|
||||
let event = create_announcement(
|
||||
@@ -160,18 +157,18 @@ async fn test_rejects_announcement_without_clone() {
|
||||
);
|
||||
|
||||
let event_msg = json!(["EVENT", event]);
|
||||
ws.send(Message::Text(event_msg.to_string()))
|
||||
ws.send(Message::Text(event_msg.to_string().into()))
|
||||
.await
|
||||
.expect("Failed to send event");
|
||||
|
||||
if let Some(Ok(Message::Text(text))) = ws.next().await {
|
||||
let response: Value = serde_json::from_str(&text).expect("Failed to parse");
|
||||
|
||||
|
||||
// Should be rejected
|
||||
assert_eq!(response[0], "OK");
|
||||
assert_eq!(response[1], event.id.to_hex());
|
||||
assert_eq!(response[2], false, "Event should be rejected");
|
||||
|
||||
|
||||
let message = response[3].as_str().unwrap();
|
||||
assert!(
|
||||
message.contains("clone") || message.contains("invalid"),
|
||||
@@ -190,9 +187,7 @@ async fn test_rejects_announcement_without_relay() {
|
||||
let relay = TestRelay::start().await;
|
||||
let keys = Keys::generate();
|
||||
|
||||
let (mut ws, _) = connect_async(relay.url())
|
||||
.await
|
||||
.expect("Failed to connect");
|
||||
let (mut ws, _) = connect_async(relay.url()).await.expect("Failed to connect");
|
||||
|
||||
// Missing relay tag
|
||||
let event = create_announcement(
|
||||
@@ -204,18 +199,18 @@ async fn test_rejects_announcement_without_relay() {
|
||||
);
|
||||
|
||||
let event_msg = json!(["EVENT", event]);
|
||||
ws.send(Message::Text(event_msg.to_string()))
|
||||
ws.send(Message::Text(event_msg.to_string().into()))
|
||||
.await
|
||||
.expect("Failed to send event");
|
||||
|
||||
if let Some(Ok(Message::Text(text))) = ws.next().await {
|
||||
let response: Value = serde_json::from_str(&text).expect("Failed to parse");
|
||||
|
||||
|
||||
// Should be rejected
|
||||
assert_eq!(response[0], "OK");
|
||||
assert_eq!(response[1], event.id.to_hex());
|
||||
assert_eq!(response[2], false, "Event should be rejected");
|
||||
|
||||
|
||||
let message = response[3].as_str().unwrap();
|
||||
assert!(
|
||||
message.contains("relays") || message.contains("invalid"),
|
||||
@@ -233,9 +228,7 @@ async fn test_rejects_announcement_for_other_service() {
|
||||
let relay = TestRelay::start().await;
|
||||
let keys = Keys::generate();
|
||||
|
||||
let (mut ws, _) = connect_async(relay.url())
|
||||
.await
|
||||
.expect("Failed to connect");
|
||||
let (mut ws, _) = connect_async(relay.url()).await.expect("Failed to connect");
|
||||
|
||||
// Lists different service
|
||||
let event = create_announcement(
|
||||
@@ -247,13 +240,13 @@ async fn test_rejects_announcement_for_other_service() {
|
||||
);
|
||||
|
||||
let event_msg = json!(["EVENT", event]);
|
||||
ws.send(Message::Text(event_msg.to_string()))
|
||||
ws.send(Message::Text(event_msg.to_string().into()))
|
||||
.await
|
||||
.expect("Failed to send event");
|
||||
|
||||
if let Some(Ok(Message::Text(text))) = ws.next().await {
|
||||
let response: Value = serde_json::from_str(&text).expect("Failed to parse");
|
||||
|
||||
|
||||
// Should be rejected
|
||||
assert_eq!(response[0], "OK");
|
||||
assert_eq!(response[1], event.id.to_hex());
|
||||
@@ -269,9 +262,7 @@ async fn test_accepts_valid_state() {
|
||||
let relay = TestRelay::start().await;
|
||||
let keys = Keys::generate();
|
||||
|
||||
let (mut ws, _) = connect_async(relay.url())
|
||||
.await
|
||||
.expect("Failed to connect");
|
||||
let (mut ws, _) = connect_async(relay.url()).await.expect("Failed to connect");
|
||||
|
||||
let event = create_state(
|
||||
&keys,
|
||||
@@ -280,13 +271,13 @@ async fn test_accepts_valid_state() {
|
||||
);
|
||||
|
||||
let event_msg = json!(["EVENT", event]);
|
||||
ws.send(Message::Text(event_msg.to_string()))
|
||||
ws.send(Message::Text(event_msg.to_string().into()))
|
||||
.await
|
||||
.expect("Failed to send event");
|
||||
|
||||
if let Some(Ok(Message::Text(text))) = ws.next().await {
|
||||
let response: Value = serde_json::from_str(&text).expect("Failed to parse");
|
||||
|
||||
|
||||
// Should be accepted
|
||||
assert_eq!(response[0], "OK");
|
||||
assert_eq!(response[1], event.id.to_hex());
|
||||
@@ -302,9 +293,7 @@ async fn test_accepts_state_with_multiple_branches() {
|
||||
let relay = TestRelay::start().await;
|
||||
let keys = Keys::generate();
|
||||
|
||||
let (mut ws, _) = connect_async(relay.url())
|
||||
.await
|
||||
.expect("Failed to connect");
|
||||
let (mut ws, _) = connect_async(relay.url()).await.expect("Failed to connect");
|
||||
|
||||
let event = create_state(
|
||||
&keys,
|
||||
@@ -317,13 +306,13 @@ async fn test_accepts_state_with_multiple_branches() {
|
||||
);
|
||||
|
||||
let event_msg = json!(["EVENT", event]);
|
||||
ws.send(Message::Text(event_msg.to_string()))
|
||||
ws.send(Message::Text(event_msg.to_string().into()))
|
||||
.await
|
||||
.expect("Failed to send event");
|
||||
|
||||
if let Some(Ok(Message::Text(text))) = ws.next().await {
|
||||
let response: Value = serde_json::from_str(&text).expect("Failed to parse");
|
||||
|
||||
|
||||
assert_eq!(response[0], "OK");
|
||||
assert_eq!(response[2], true, "State event should be accepted");
|
||||
} else {
|
||||
@@ -337,9 +326,7 @@ async fn test_rejects_state_without_identifier() {
|
||||
let relay = TestRelay::start().await;
|
||||
let keys = Keys::generate();
|
||||
|
||||
let (mut ws, _) = connect_async(relay.url())
|
||||
.await
|
||||
.expect("Failed to connect");
|
||||
let (mut ws, _) = connect_async(relay.url()).await.expect("Failed to connect");
|
||||
|
||||
// Create state without identifier
|
||||
let event = EventBuilder::new(Kind::from(KIND_REPOSITORY_STATE), "")
|
||||
@@ -347,18 +334,18 @@ async fn test_rejects_state_without_identifier() {
|
||||
.expect("Failed to sign event");
|
||||
|
||||
let event_msg = json!(["EVENT", event]);
|
||||
ws.send(Message::Text(event_msg.to_string()))
|
||||
ws.send(Message::Text(event_msg.to_string().into()))
|
||||
.await
|
||||
.expect("Failed to send event");
|
||||
|
||||
if let Some(Ok(Message::Text(text))) = ws.next().await {
|
||||
let response: Value = serde_json::from_str(&text).expect("Failed to parse");
|
||||
|
||||
|
||||
// Should be rejected
|
||||
assert_eq!(response[0], "OK");
|
||||
assert_eq!(response[1], event.id.to_hex());
|
||||
assert_eq!(response[2], false, "Event should be rejected");
|
||||
|
||||
|
||||
let message = response[3].as_str().unwrap();
|
||||
assert!(
|
||||
message.contains("identifier") || message.contains("invalid"),
|
||||
@@ -376,21 +363,22 @@ async fn test_query_announcements() {
|
||||
let relay = TestRelay::start().await;
|
||||
let keys = Keys::generate();
|
||||
|
||||
let (mut ws, _) = connect_async(relay.url())
|
||||
.await
|
||||
.expect("Failed to connect");
|
||||
let (mut ws, _) = connect_async(relay.url()).await.expect("Failed to connect");
|
||||
|
||||
// Send an announcement
|
||||
let event = create_announcement(
|
||||
&keys,
|
||||
&relay.domain(),
|
||||
"query-test-repo",
|
||||
vec![format!("https://{}/alice/query-test-repo.git", relay.domain())],
|
||||
vec![format!(
|
||||
"https://{}/alice/query-test-repo.git",
|
||||
relay.domain()
|
||||
)],
|
||||
vec![format!("wss://{}", relay.domain())],
|
||||
);
|
||||
|
||||
let event_msg = json!(["EVENT", event]);
|
||||
ws.send(Message::Text(event_msg.to_string()))
|
||||
ws.send(Message::Text(event_msg.to_string().into()))
|
||||
.await
|
||||
.expect("Failed to send event");
|
||||
|
||||
@@ -409,7 +397,7 @@ async fn test_query_announcements() {
|
||||
}
|
||||
]);
|
||||
|
||||
ws.send(Message::Text(req.to_string()))
|
||||
ws.send(Message::Text(req.to_string().into()))
|
||||
.await
|
||||
.expect("Failed to send REQ");
|
||||
|
||||
@@ -420,7 +408,7 @@ async fn test_query_announcements() {
|
||||
for _ in 0..10 {
|
||||
if let Some(Ok(Message::Text(text))) = ws.next().await {
|
||||
let response: Value = serde_json::from_str(&text).expect("Failed to parse");
|
||||
|
||||
|
||||
if response[0] == "EVENT" {
|
||||
assert_eq!(response[1], "test-sub");
|
||||
found_event = true;
|
||||
@@ -442,9 +430,7 @@ async fn test_query_states() {
|
||||
let relay = TestRelay::start().await;
|
||||
let keys = Keys::generate();
|
||||
|
||||
let (mut ws, _) = connect_async(relay.url())
|
||||
.await
|
||||
.expect("Failed to connect");
|
||||
let (mut ws, _) = connect_async(relay.url()).await.expect("Failed to connect");
|
||||
|
||||
// Send a state event
|
||||
let event = create_state(
|
||||
@@ -454,7 +440,7 @@ async fn test_query_states() {
|
||||
);
|
||||
|
||||
let event_msg = json!(["EVENT", event]);
|
||||
ws.send(Message::Text(event_msg.to_string()))
|
||||
ws.send(Message::Text(event_msg.to_string().into()))
|
||||
.await
|
||||
.expect("Failed to send event");
|
||||
|
||||
@@ -473,7 +459,7 @@ async fn test_query_states() {
|
||||
}
|
||||
]);
|
||||
|
||||
ws.send(Message::Text(req.to_string()))
|
||||
ws.send(Message::Text(req.to_string().into()))
|
||||
.await
|
||||
.expect("Failed to send REQ");
|
||||
|
||||
@@ -484,7 +470,7 @@ async fn test_query_states() {
|
||||
for _ in 0..10 {
|
||||
if let Some(Ok(Message::Text(text))) = ws.next().await {
|
||||
let response: Value = serde_json::from_str(&text).expect("Failed to parse");
|
||||
|
||||
|
||||
if response[0] == "EVENT" {
|
||||
assert_eq!(response[1], "test-sub");
|
||||
found_event = true;
|
||||
@@ -506,21 +492,22 @@ async fn test_duplicate_announcement() {
|
||||
let relay = TestRelay::start().await;
|
||||
let keys = Keys::generate();
|
||||
|
||||
let (mut ws, _) = connect_async(relay.url())
|
||||
.await
|
||||
.expect("Failed to connect");
|
||||
let (mut ws, _) = connect_async(relay.url()).await.expect("Failed to connect");
|
||||
|
||||
let event = create_announcement(
|
||||
&keys,
|
||||
&relay.domain(),
|
||||
"duplicate-test",
|
||||
vec![format!("https://{}/alice/duplicate-test.git", relay.domain())],
|
||||
vec![format!(
|
||||
"https://{}/alice/duplicate-test.git",
|
||||
relay.domain()
|
||||
)],
|
||||
vec![format!("wss://{}", relay.domain())],
|
||||
);
|
||||
|
||||
// Send first time
|
||||
let event_msg = json!(["EVENT", event]);
|
||||
ws.send(Message::Text(event_msg.to_string()))
|
||||
ws.send(Message::Text(event_msg.to_string().into()))
|
||||
.await
|
||||
.expect("Failed to send event");
|
||||
|
||||
@@ -531,14 +518,14 @@ async fn test_duplicate_announcement() {
|
||||
|
||||
// Send second time (duplicate)
|
||||
let event_msg = json!(["EVENT", event]);
|
||||
ws.send(Message::Text(event_msg.to_string()))
|
||||
ws.send(Message::Text(event_msg.to_string().into()))
|
||||
.await
|
||||
.expect("Failed to send event");
|
||||
|
||||
if let Some(Ok(Message::Text(text))) = ws.next().await {
|
||||
let response2: Value = serde_json::from_str(&text).expect("Failed to parse");
|
||||
assert_eq!(response2[2], true, "Duplicate should be acknowledged");
|
||||
|
||||
|
||||
let message = response2[3].as_str().unwrap();
|
||||
assert!(
|
||||
message.contains("duplicate") || message.is_empty(),
|
||||
|
||||
Reference in New Issue
Block a user