More docs and expand key_value_store example

This commit is contained in:
David Pedersen
2021-06-08 12:43:16 +02:00
parent d7605d3184
commit 1f8b39f05d
13 changed files with 920 additions and 431 deletions
+88 -18
View File
@@ -1,47 +1,70 @@
//! Simple in-memory key/value store showing features of tower-web.
//!
//! Run with:
//!
//! ```not_rust
//! RUST_LOG=tower_http=debug,key_value_store=trace cargo run --example key_value_store
//! ```
use bytes::Bytes;
use http::{Request, StatusCode};
use hyper::Server;
use std::{
borrow::Cow,
collections::HashMap,
net::SocketAddr,
sync::{Arc, Mutex},
sync::{Arc, RwLock},
time::Duration,
};
use tower::{make::Shared, ServiceBuilder};
use tower::{make::Shared, BoxError, ServiceBuilder};
use tower_http::{
add_extension::AddExtensionLayer, compression::CompressionLayer, trace::TraceLayer,
add_extension::AddExtensionLayer, auth::RequireAuthorizationLayer,
compression::CompressionLayer, trace::TraceLayer,
};
use tower_web::{
body::Body,
body::{Body, BoxBody},
extract::{BytesMaxLength, Extension, UrlParams},
prelude::*,
response::IntoResponse,
routing::BoxRoute,
};
#[tokio::main]
async fn main() {
tracing_subscriber::fmt::init();
// build our application with some routes
// Build our application by composing routes
let app = route(
"/:key",
get(kv_get.layer(CompressionLayer::new())).post(kv_set),
);
// Add compression to `kv_get`
get(kv_get.layer(CompressionLayer::new()))
// But don't compress `kv_set`
.post(kv_set),
)
.route("/keys", get(list_keys))
// Nest our admin routes under `/admin`
.nest("/admin", admin_routes())
// Add middleware to all routes
.layer(
ServiceBuilder::new()
.load_shed()
.concurrency_limit(1024)
.timeout(Duration::from_secs(10))
.layer(TraceLayer::new_for_http())
.layer(AddExtensionLayer::new(SharedState::default()))
.into_inner(),
)
// Handle errors from middleware
.handle_error(handle_error);
// add some middleware
let app = ServiceBuilder::new()
.timeout(Duration::from_secs(10))
.layer(TraceLayer::new_for_http())
.layer(AddExtensionLayer::new(SharedState::default()))
.service(app);
// run it with hyper
// Run our app with hyper
let addr = SocketAddr::from(([127, 0, 0, 1], 3000));
tracing::debug!("listening on {}", addr);
let server = Server::bind(&addr).serve(Shared::new(app));
server.await.unwrap();
}
type SharedState = Arc<Mutex<State>>;
type SharedState = Arc<RwLock<State>>;
#[derive(Default)]
struct State {
@@ -53,7 +76,7 @@ async fn kv_get(
UrlParams((key,)): UrlParams<(String,)>,
Extension(state): Extension<SharedState>,
) -> Result<Bytes, StatusCode> {
let db = &state.lock().unwrap().db;
let db = &state.read().unwrap().db;
if let Some(value) = db.get(&key) {
Ok(value.clone())
@@ -68,5 +91,52 @@ async fn kv_set(
BytesMaxLength(value): BytesMaxLength<{ 1024 * 5_000 }>, // ~5mb
Extension(state): Extension<SharedState>,
) {
state.lock().unwrap().db.insert(key, value);
state.write().unwrap().db.insert(key, value);
}
async fn list_keys(_req: Request<Body>, Extension(state): Extension<SharedState>) -> String {
let db = &state.read().unwrap().db;
db.keys()
.map(|key| key.to_string())
.collect::<Vec<String>>()
.join("\n")
}
fn admin_routes() -> BoxRoute<BoxBody> {
async fn delete_all_keys(_req: Request<Body>, Extension(state): Extension<SharedState>) {
state.write().unwrap().db.clear();
}
async fn remove_key(
_req: Request<Body>,
UrlParams((key,)): UrlParams<(String,)>,
Extension(state): Extension<SharedState>,
) {
state.write().unwrap().db.remove(&key);
}
route("/keys", delete(delete_all_keys))
.route("/key/:key", delete(remove_key))
// Require beare auth for all admin routes
.layer(RequireAuthorizationLayer::bearer("secret-token"))
.boxed()
}
fn handle_error(error: BoxError) -> impl IntoResponse {
if error.is::<tower::timeout::error::Elapsed>() {
return (StatusCode::REQUEST_TIMEOUT, Cow::from("request timed out"));
}
if error.is::<tower::load_shed::error::Overloaded>() {
return (
StatusCode::SERVICE_UNAVAILABLE,
Cow::from("service is overloaded, try again later"),
);
}
(
StatusCode::INTERNAL_SERVER_ERROR,
Cow::from(format!("Unhandled internal error: {}", error)),
)
}
+1 -1
View File
@@ -3,7 +3,7 @@ use hyper::Server;
use std::net::SocketAddr;
use tower::make::Shared;
use tower_http::{services::ServeDir, trace::TraceLayer};
use tower_web::{prelude::*, ServiceExt};
use tower_web::{prelude::*, service::ServiceExt};
#[tokio::main]
async fn main() {