1- use axum:: extract:: State ;
1+ use axum:: extract:: State ;
22use axum:: http:: { HeaderMap , StatusCode } ;
33use axum:: response:: { IntoResponse , Response , Sse } ;
44use axum:: routing:: { get, post} ;
@@ -22,8 +22,8 @@ use crate::media_guard::{
2222 unsupported_multimodal_error_message,
2323} ;
2424use crate :: project:: {
25- current_environment , project_id_from_key, project_key_from_id, read_project_env,
26- registry_dir_path , validate_adapter_token, ProjectRegistry , PROJECT_ENV_FILENAME ,
25+ project_id_from_key, project_key_from_id, read_project_env, registry_dir_path ,
26+ validate_adapter_token, ProjectRegistry , PROJECT_ENV_FILENAME ,
2727} ;
2828use crate :: state:: { now_ts, StateStore } ;
2929use crate :: upstream:: {
@@ -144,7 +144,7 @@ async fn responses(State(state): State<AppState>, headers: HeaderMap, body: Stri
144144 if let Err ( response) = authorize_adapter ( & projects, & headers) {
145145 return * response;
146146 }
147- match projects . get ( & project_id) {
147+ match project_runtime ( & projects , & project_id) {
148148 Some ( r) => r. clone ( ) ,
149149 None => {
150150 return error_response (
@@ -222,7 +222,7 @@ async fn complete_response(
222222 & message,
223223 ) ;
224224 }
225- let permit = match capacity. clone ( ) . try_acquire_owned ( ) {
225+ let _permit = match capacity. clone ( ) . try_acquire_owned ( ) {
226226 Ok ( permit) => permit,
227227 Err ( _) => {
228228 return error_response (
@@ -233,7 +233,7 @@ async fn complete_response(
233233 }
234234 } ;
235235 let upstream = runtime. client . chat ( payload) . await ;
236- drop ( permit ) ;
236+ drop ( _permit ) ;
237237 let upstream = match upstream {
238238 Ok ( value) => value,
239239 Err ( error) => {
@@ -346,10 +346,18 @@ async fn stream_response(
346346 ) ;
347347 if let Err ( error) = assembler. start ( ) {
348348 tracing:: error!( error = %error, "failed to emit initial stream lifecycle events" ) ;
349+
350+ let response =
351+ responses_failed_value ( & body, & model_alias, "internal_error" , & error. to_string ( ) ) ;
352+ let event = json ! ( { "type" : "response.failed" , "response" : response} ) ;
353+ let _ = tx. send ( Ok ( axum:: response:: sse:: Event :: default ( )
354+ . event ( "response.failed" )
355+ . data ( event. to_string ( ) ) ) ) ;
349356 let _ = tx. send ( Ok ( axum:: response:: sse:: Event :: default ( ) . data ( "[DONE]" ) ) ) ;
357+
350358 return ;
351359 }
352- let permit = match capacity. clone ( ) . try_acquire_owned ( ) {
360+ let _permit = match capacity. clone ( ) . try_acquire_owned ( ) {
353361 Ok ( permit) => permit,
354362 Err ( _) => {
355363 if let Err ( error) =
@@ -362,7 +370,7 @@ async fn stream_response(
362370 }
363371 } ;
364372 let upstream = runtime_for_task. client . chat_stream ( payload) . await ;
365- drop ( permit ) ;
373+
366374 match upstream {
367375 Ok ( mut stream) => {
368376 let mut buffer = String :: new ( ) ;
@@ -576,10 +584,7 @@ fn previous_response(
576584 Ok ( Some ( previous) )
577585}
578586
579- fn authorize_adapter (
580- projects : & HashMap < String , ProjectRuntime > ,
581- headers : & HeaderMap ,
582- ) -> Result < ( ) , Box < Response > > {
587+ fn adapter_bearer_token ( headers : & HeaderMap ) -> Result < & str , Box < Response > > {
583588 let auth_header = headers
584589 . get ( "authorization" )
585590 . and_then ( |v| v. to_str ( ) . ok ( ) )
@@ -590,19 +595,48 @@ fn authorize_adapter(
590595 "Missing Authorization header. Provide a valid Bearer token." ,
591596 ) )
592597 } ) ?;
593- let raw_token = auth_header. strip_prefix ( "Bearer " ) . ok_or_else ( || {
598+ auth_header. strip_prefix ( "Bearer " ) . ok_or_else ( || {
594599 Box :: new ( error_response (
595600 StatusCode :: UNAUTHORIZED ,
596601 "unauthorized" ,
597602 "Invalid Authorization format. Expected 'Bearer <token>'." ,
598603 ) )
599- } ) ?;
600- for runtime in projects. values ( ) {
601- if let Some ( ref local_token) = runtime. config . local_token {
602- if !local_token. is_empty ( ) && validate_adapter_token ( raw_token, local_token) {
603- return Ok ( ( ) ) ;
604- }
605- }
604+ } )
605+ }
606+
607+ fn project_runtime < ' a > (
608+ projects : & ' a HashMap < String , ProjectRuntime > ,
609+ project_id : & str ,
610+ ) -> Option < & ' a ProjectRuntime > {
611+ projects
612+ . get ( project_id)
613+ . or_else ( || projects. get ( project_key_from_id ( project_id) ) )
614+ . or_else ( || {
615+ let canonical_id = project_id_from_key ( project_id) ;
616+ projects. get ( & canonical_id)
617+ } )
618+ }
619+
620+ fn runtime_accepts_token ( runtime : & ProjectRuntime , raw_token : & str ) -> bool {
621+ runtime
622+ . config
623+ . local_token
624+ . as_ref ( )
625+ . is_some_and ( |local_token| {
626+ !local_token. is_empty ( ) && validate_adapter_token ( raw_token, local_token)
627+ } )
628+ }
629+
630+ fn authorize_adapter (
631+ projects : & HashMap < String , ProjectRuntime > ,
632+ headers : & HeaderMap ,
633+ ) -> Result < ( ) , Box < Response > > {
634+ let raw_token = adapter_bearer_token ( headers) ?;
635+ if projects
636+ . values ( )
637+ . any ( |runtime| runtime_accepts_token ( runtime, raw_token) )
638+ {
639+ return Ok ( ( ) ) ;
606640 }
607641 Err ( Box :: new ( error_response (
608642 StatusCode :: UNAUTHORIZED ,
@@ -611,10 +645,7 @@ fn authorize_adapter(
611645 ) ) )
612646}
613647
614- async fn admin_refresh (
615- State ( state) : State < AppState > ,
616- headers : HeaderMap ,
617- ) -> Response {
648+ async fn admin_refresh ( State ( state) : State < AppState > , headers : HeaderMap ) -> Response {
618649 // Auth: accept any valid adapter bearer token.
619650 let auth_ok = {
620651 let projects = state. projects . read ( ) . unwrap ( ) ;
@@ -657,16 +688,20 @@ async fn admin_refresh(
657688 continue ;
658689 }
659690 } ;
660- let env = current_environment ( ) ;
661- let config = match Config :: from_sources ( & project_env, & env, state. config_overrides . clone ( ) ) {
691+ let env = HashMap :: new ( ) ;
692+ let config = match Config :: from_sources ( & project_env, & env, state. config_overrides . clone ( ) )
693+ {
662694 Ok ( c) => c,
663695 Err ( e) => {
664696 tracing:: warn!( "refresh: bad config for {project_id}: {e}" ) ;
665697 continue ;
666698 }
667699 } ;
668700 let state_db_path = root. join ( & config. state_db ) ;
669- let store = match StateStore :: new ( state_db_path. display ( ) . to_string ( ) , config. state_ttl_seconds ) {
701+ let store = match StateStore :: new (
702+ state_db_path. display ( ) . to_string ( ) ,
703+ config. state_ttl_seconds ,
704+ ) {
670705 Ok ( s) => s,
671706 Err ( e) => {
672707 tracing:: warn!( "refresh: cannot create state for {project_id}: {e}" ) ;
@@ -684,14 +719,20 @@ async fn admin_refresh(
684719 continue ;
685720 }
686721 } ;
687- projects. insert ( project_id. clone ( ) , ProjectRuntime { config, client, state : store } ) ;
722+ projects. insert (
723+ project_id. clone ( ) ,
724+ ProjectRuntime {
725+ config,
726+ client,
727+ state : store,
728+ } ,
729+ ) ;
688730 added. push ( project_id. clone ( ) ) ;
689731 }
690732
691733 Json ( json ! ( { "status" : "ok" , "added" : added, "already_loaded" : skipped} ) ) . into_response ( )
692734}
693735
694-
695736fn parse_routed_model ( model : & str ) -> Result < ( String , String ) , & ' static str > {
696737 let Some ( rest) = model. strip_prefix ( "opencode_adapter/" ) else {
697738 return Err ( "model must use opencode_adapter/<project_key>/<real_model>. Run 'codex-opencode-adapter init' to refresh agent templates." ) ;
@@ -863,6 +904,3 @@ fn upstream_error(error: UpstreamError) -> Response {
863904 }
864905 }
865906}
866-
867-
868-
0 commit comments