| author | |
| committer | |
| log | d9b295aa93ef00565b91d9caa3ec4875b94596a4 |
| tree | b83e1578da850692b072ab5fa6e213c93628fed3 |
| parent | fed95315bf3a060a272a75be193773cde19900f9 |
| signature | Signed by SSH key SHA256:52mNGHRsVFBDED9IAX5pe+LRWUefqTbxEReunq21QvU |
Keep upstream streaming independent of the dashboard request timeout.
Assisted-by: gpt-63 files changed, 39 insertions(+), 16 deletions(-)
dashboard/src/main.rs+18| ... | @@ -91,6 +91,7 @@ struct App { | ... | @@ -91,6 +91,7 @@ struct App { |
| 91 | relay: relay::Broker, | 91 | relay: relay::Broker, |
| 92 | shale: shale::Backend, | 92 | shale: shale::Backend, |
| 93 | http: reqwest::Client, | 93 | http: reqwest::Client, |
| 94 | shale_http: reqwest::Client, | ||
| 94 | internal: Option<(url::Url, reqwest::Client)>, | 95 | internal: Option<(url::Url, reqwest::Client)>, |
| 95 | cache: cache::Cache, | 96 | cache: cache::Cache, |
| 96 | nomad_slots: Semaphore, | 97 | nomad_slots: Semaphore, |
| ... | @@ -382,6 +383,22 @@ async fn main() -> std::result::Result<(), Box<dyn std::error::Error>> { | ... | @@ -382,6 +383,22 @@ async fn main() -> std::result::Result<(), Box<dyn std::error::Error>> { |
| 382 | }, | 383 | }, |
| 383 | ) | 384 | ) |
| 384 | .transpose()?; | 385 | .transpose()?; |
| 386 | let mut shale_http = reqwest::Client::builder() | ||
| 387 | .connect_timeout(Duration::from_secs(5)) | ||
| 388 | .read_timeout(Duration::from_secs(30)) | ||
| 389 | .redirect(reqwest::redirect::Policy::none()) | ||
| 390 | .retry(reqwest::retry::never()); | ||
| 391 | for certificate in &certificates { | ||
| 392 | shale_http = shale_http.add_root_certificate(certificate.clone()); | ||
| 393 | } | ||
| 394 | if internal.is_some() { | ||
| 395 | let mut token = axum::http::HeaderValue::from_str(proof.as_ref().unwrap())?; | ||
| 396 | token.set_sensitive(true); | ||
| 397 | let mut headers = HeaderMap::new(); | ||
| 398 | headers.insert("Studio-Proxy-Token", token); | ||
| 399 | shale_http = shale_http.default_headers(headers); | ||
| 400 | } | ||
| 401 | let shale_http = shale_http.build()?; | ||
| 385 | let (live, _) = watch::channel(Bytes::new()); | 402 | let (live, _) = watch::channel(Bytes::new()); |
| 386 | let index = std::env::var("STUDIO_INDEX_POOL").ok().map(|pool| { | 403 | let index = std::env::var("STUDIO_INDEX_POOL").ok().map(|pool| { |
| 387 | Arc::new(index::Index::new( | 404 | Arc::new(index::Index::new( |
| ... | @@ -459,6 +476,7 @@ async fn main() -> std::result::Result<(), Box<dyn std::error::Error>> { | ... | @@ -459,6 +476,7 @@ async fn main() -> std::result::Result<(), Box<dyn std::error::Error>> { |
| 459 | ) | 476 | ) |
| 460 | .map_err(|error| std::io::Error::other(error.message))?, | 477 | .map_err(|error| std::io::Error::other(error.message))?, |
| 461 | http: client().build()?, | 478 | http: client().build()?, |
| 479 | shale_http, | ||
| 462 | internal, | 480 | internal, |
| 463 | cache: cache::Cache::default(), | 481 | cache: cache::Cache::default(), |
| 464 | nomad_slots: Semaphore::new(4), | 482 | nomad_slots: Semaphore::new(4), |
dashboard/src/shale_page.rs+10-11| ... | @@ -182,6 +182,13 @@ pub async fn proxy(app: Arc<App>, request: Request) -> Result<Response> { | ... | @@ -182,6 +182,13 @@ pub async fn proxy(app: Arc<App>, request: Request) -> Result<Response> { |
| 182 | ) | 182 | ) |
| 183 | .into_response()); | 183 | .into_response()); |
| 184 | } | 184 | } |
| 185 | let target = match &app.internal { | ||
| 186 | Some((base, _)) => format!("{base}services/{service}{uri}"), | ||
| 187 | None => format!("http://{upstream}{uri}"), | ||
| 188 | }; | ||
| 189 | if app.internal.is_some() { | ||
| 190 | parts.headers.remove("host"); | ||
| 191 | } | ||
| 185 | strip_hop_headers(&mut parts.headers); | 192 | strip_hop_headers(&mut parts.headers); |
| 186 | for name in [ | 193 | for name in [ |
| 187 | "studio-shale-upstream", | 194 | "studio-shale-upstream", |
| ... | @@ -199,17 +206,9 @@ pub async fn proxy(app: Arc<App>, request: Request) -> Result<Response> { | ... | @@ -199,17 +206,9 @@ pub async fn proxy(app: Arc<App>, request: Request) -> Result<Response> { |
| 199 | parts | 206 | parts |
| 200 | .headers | 207 | .headers |
| 201 | .insert("accept-encoding", "identity".parse().unwrap()); | 208 | .insert("accept-encoding", "identity".parse().unwrap()); |
| 202 | static HTTP: std::sync::LazyLock<reqwest::Client> = std::sync::LazyLock::new(|| { | 209 | let mut upstream_response = app |
| 203 | reqwest::Client::builder() | 210 | .shale_http |
| 204 | .connect_timeout(Duration::from_secs(5)) | 211 | .request(parts.method.clone(), target) |
| 205 | .read_timeout(Duration::from_secs(30)) | ||
| 206 | .redirect(reqwest::redirect::Policy::none()) | ||
| 207 | .retry(reqwest::retry::never()) | ||
| 208 | .build() | ||
| 209 | .unwrap() | ||
| 210 | }); | ||
| 211 | let mut upstream_response = HTTP | ||
| 212 | .request(parts.method.clone(), format!("http://{upstream}{uri}")) | ||
| 213 | .headers(parts.headers) | 212 | .headers(parts.headers) |
| 214 | .body(reqwest::Body::wrap_stream(body.into_data_stream())) | 213 | .body(reqwest::Body::wrap_stream(body.into_data_stream())) |
| 215 | .send() | 214 | .send() |
tools/dashboard-shale-page-test.py+11-5| ... | @@ -52,7 +52,7 @@ def main(): | ... | @@ -52,7 +52,7 @@ def main(): |
| 52 | protocol_version = 'HTTP/1.1' | 52 | protocol_version = 'HTTP/1.1' |
| 53 | 53 | ||
| 54 | def do_GET(self): | 54 | def do_GET(self): |
| 55 | observations.append((self.command, self.path, dict(self.headers))) | 55 | observations.append((self.command, self.path, {key.lower(): value for key, value in self.headers.items()})) |
| 56 | if self.path == '/binary': | 56 | if self.path == '/binary': |
| 57 | return self.respond(blob, 'application/octet-stream') | 57 | return self.respond(blob, 'application/octet-stream') |
| 58 | if self.path == '/large': | 58 | if self.path == '/large': |
| ... | @@ -69,7 +69,7 @@ def main(): | ... | @@ -69,7 +69,7 @@ def main(): |
| 69 | 69 | ||
| 70 | def do_POST(self): | 70 | def do_POST(self): |
| 71 | body = self.rfile.read(int(self.headers.get('Content-Length', 0))) | 71 | body = self.rfile.read(int(self.headers.get('Content-Length', 0))) |
| 72 | observations.append((self.command, self.path, dict(self.headers), body)) | 72 | observations.append((self.command, self.path, {key.lower(): value for key, value in self.headers.items()}, body)) |
| 73 | self.respond(body, 'application/octet-stream', 201) | 73 | self.respond(body, 'application/octet-stream', 201) |
| 74 | 74 | ||
| 75 | def respond(self, body, content_type, status=200, headers=()): | 75 | def respond(self, body, content_type, status=200, headers=()): |
| ... | @@ -95,13 +95,15 @@ def main(): | ... | @@ -95,13 +95,15 @@ def main(): |
| 95 | stack.callback(fixture.server_close) | 95 | stack.callback(fixture.server_close) |
| 96 | stack.callback(fixture.shutdown) | 96 | stack.callback(fixture.shutdown) |
| 97 | threading.Thread(target=fixture.serve_forever, daemon=True).start() | 97 | threading.Thread(target=fixture.serve_forever, daemon=True).start() |
| 98 | dashboard_port, gateway_port = port(), port() | 98 | dashboard_port, gateway_port, internal_port = port(), port(), port() |
| 99 | proof = 'a' * 64 | 99 | proof = 'a' * 64 |
| 100 | token = root / 'proxy.token' | 100 | token = root / 'proxy.token' |
| 101 | token.write_text(proof) | 101 | token.write_text(proof) |
| 102 | certificate, gateway_key = root / 'gateway.pem', root / 'gateway.key' | ||
| 103 | subprocess.run(['openssl', 'req', '-x509', '-newkey', 'rsa:2048', '-nodes', '-days', '1', '-subj', '/CN=localhost', '-addext', 'subjectAltName=DNS:localhost', '-keyout', str(gateway_key), '-out', str(certificate)], check=True, capture_output=True) | ||
| 102 | data = root / 'data' | 104 | data = root / 'data' |
| 103 | environment = {**os.environ, 'PORT': str(dashboard_port), 'STUDIO_DOMAIN': 'studio.test', 'STUDIO_DATA_DIR': str(data), 'STUDIO_PROXY_TOKEN_FILE': str(token), 'STUDIO_REPO': str(repo), 'STUDIO_WEB_DIR': str(repo / 'dashboard/dist')} | 105 | environment = {**os.environ, 'PORT': str(dashboard_port), 'STUDIO_DOMAIN': 'studio.test', 'STUDIO_DATA_DIR': str(data), 'STUDIO_PROXY_TOKEN_FILE': str(token), 'STUDIO_INTERNAL_URL': f'https://localhost:{internal_port}', 'STUDIO_CA_BUNDLE': str(certificate), 'STUDIO_REPO': str(repo), 'STUDIO_WEB_DIR': str(repo / 'dashboard/dist')} |
| 104 | for key in ['STUDIO_INTERNAL_URL', 'STUDIO_INDEX_POOL', 'STUDIO_AUTH_REQUIRED', 'STUDIO_YT_STATE']: | 106 | for key in ['STUDIO_INDEX_POOL', 'STUDIO_AUTH_REQUIRED', 'STUDIO_YT_STATE']: |
| 105 | environment.pop(key, None) | 107 | environment.pop(key, None) |
| 106 | log = stack.enter_context((root / 'dashboard.log').open('w')) | 108 | log = stack.enter_context((root / 'dashboard.log').open('w')) |
| 107 | dashboard = subprocess.Popen([str(args.binary)], env=environment, stdout=log, stderr=log) | 109 | dashboard = subprocess.Popen([str(args.binary)], env=environment, stdout=log, stderr=log) |
| ... | @@ -149,6 +151,10 @@ def main(): | ... | @@ -149,6 +151,10 @@ def main(): |
| 149 | 151 | ||
| 150 | preview_port = port() | 152 | preview_port = port() |
| 151 | config = '{\n admin off\n auto_https off\n}\n' + local_site('shale.studio.test', gateway_port) + '\n' + local_site('shale-preview-12345678.studio.test', preview_port) | 153 | config = '{\n admin off\n auto_https off\n}\n' + local_site('shale.studio.test', gateway_port) + '\n' + local_site('shale-preview-12345678.studio.test', preview_port) |
| 154 | config += f'\nhttps://localhost:{internal_port} {{\n tls {certificate} {gateway_key}\n @trusted header Studio-Proxy-Token {proof}\n handle @trusted {{\n request_header -Studio-Proxy-Token\n' | ||
| 155 | for service in ['shale', 'shale-preview-12345678']: | ||
| 156 | config += f' handle_path /services/{service}/* {{\n reverse_proxy 127.0.0.1:{fixture.server_port}\n }}\n' | ||
| 157 | config += ' }\n handle {\n respond 403\n }\n}\n' | ||
| 152 | config_path = root / 'Caddyfile' | 158 | config_path = root / 'Caddyfile' |
| 153 | config_path.write_text(config) | 159 | config_path.write_text(config) |
| 154 | subprocess.run([str(args.caddy), 'validate', '--config', str(config_path), '--adapter', 'caddyfile'], check=True, capture_output=True, env=caddy_env) | 160 | subprocess.run([str(args.caddy), 'validate', '--config', str(config_path), '--adapter', 'caddyfile'], check=True, capture_output=True, env=caddy_env) |