Skip to content

Commit 799cc9f

Browse files
authored
feat: bound global admission and expose running backend traces (#11560)
feat: bound backend admission and expose running traces Add process-wide backend execution admission without blocking UI or administrative HTTP work. Represent backend operations while they are in flight, surface running traces with immediate log links, and tie streaming admission leases to the gRPC receive lifecycle. Assisted-by: OpenAI Codex: GPT-5 Signed-off-by: Richard Palethorpe <io@richiejp.com>
1 parent 2ae7b45 commit 799cc9f

47 files changed

Lines changed: 852 additions & 78 deletions

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

core/application/application.go

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -25,6 +25,7 @@ import (
2525
"github.com/mudler/LocalAI/core/services/voiceprofile"
2626
"github.com/mudler/LocalAI/core/services/voicerecognition"
2727
"github.com/mudler/LocalAI/core/templates"
28+
"github.com/mudler/LocalAI/core/trace"
2829
pkggrpc "github.com/mudler/LocalAI/pkg/grpc"
2930
localaitools "github.com/mudler/LocalAI/pkg/mcp/localaitools"
3031
localaiInproc "github.com/mudler/LocalAI/pkg/mcp/localaitools/inproc"
@@ -119,6 +120,8 @@ func (a *Application) Ready() bool { return a.startupComplete.Load() }
119120
func (a *Application) markStartupComplete() { a.startupComplete.Store(true) }
120121

121122
func newApplication(appConfig *config.ApplicationConfig) *Application {
123+
corebackend.ConfigureGlobalBackendAdmission(appConfig.MaxConcurrentBackendRequests)
124+
trace.ConfigureBackendTraceMaxInFlight(appConfig.MaxConcurrentBackendRequests)
122125
ml := model.NewModelLoader(appConfig.SystemState)
123126

124127
// Apply the per-model load-failure cooldown (0 disables). Set here rather
@@ -134,7 +137,7 @@ func newApplication(appConfig *config.ApplicationConfig) *Application {
134137
// Record a model_load backend trace for every real backend load, so the
135138
// Traces UI shows which backend runtime served each model and how long
136139
// the load took. Load failures are traced by the modality wrappers.
137-
ml.SetLoadObserver(corebackend.ModelLoadTraceObserver(appConfig))
140+
ml.SetLoadLifecycleObserver(corebackend.ModelLoadTraceObserver(appConfig))
138141

139142
app := &Application{
140143
backendLoader: config.NewModelConfigLoader(

core/backend/audio_transform.go

Lines changed: 20 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -91,12 +91,20 @@ func ModelAudioTransform(
9191
return AudioTransformOutputs{}, nil, fmt.Errorf("persist reference: %w", err)
9292
}
9393
}
94+
release, err := AcquireGlobalBackendSlot()
95+
if err != nil {
96+
return AudioTransformOutputs{}, nil, err
97+
}
98+
defer release()
9499

95100
var startTime time.Time
101+
var traceID string
96102
if appConfig.EnableTracing {
97103
trace.InitBackendTracingIfEnabled(appConfig.TracingMaxItems, appConfig.TracingMaxBodyBytes)
98104
startTime = time.Now()
105+
traceID = trace.BeginBackendTrace(trace.BackendTrace{Timestamp: startTime, Type: trace.BackendTraceAudioTransform, ModelName: modelConfig.Name, Backend: modelConfig.Backend, Summary: trace.TruncateString(filepath.Base(audioPath), 200)})
99106
}
107+
defer trace.CancelBackendTrace(traceID)
100108

101109
res, err := transformModel.AudioTransform(ctx, &proto.AudioTransformRequest{
102110
ModelIdentity: modelConfig.Model,
@@ -126,6 +134,7 @@ func ModelAudioTransform(
126134
}
127135
}
128136
trace.RecordBackendTrace(trace.BackendTrace{
137+
ID: traceID,
129138
Timestamp: startTime,
130139
Duration: time.Since(startTime),
131140
Type: trace.BackendTraceAudioTransform,
@@ -198,7 +207,17 @@ func ModelAudioTransformStream(
198207
if transformModel == nil {
199208
return nil, fmt.Errorf("could not load audio-transform model %q", modelConfig.Model)
200209
}
201-
return transformModel.AudioTransformStream(ctx)
210+
release, err := AcquireGlobalBackendSlot()
211+
if err != nil {
212+
return nil, err
213+
}
214+
stream, err := transformModel.AudioTransformStream(ctx)
215+
if err != nil {
216+
release()
217+
return nil, err
218+
}
219+
stream.AddCleanup(release)
220+
return stream, nil
202221
}
203222

204223
// persistAudioInput copies a transient input file (typically a multipart

core/backend/depth.go

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -33,12 +33,20 @@ func Depth(
3333
if depthModel == nil {
3434
return nil, fmt.Errorf("could not load depth model")
3535
}
36+
release, err := AcquireGlobalBackendSlot()
37+
if err != nil {
38+
return nil, err
39+
}
40+
defer release()
3641

3742
var startTime time.Time
43+
var traceID string
3844
if appConfig.EnableTracing {
3945
trace.InitBackendTracingIfEnabled(appConfig.TracingMaxItems, appConfig.TracingMaxBodyBytes)
4046
startTime = time.Now()
47+
traceID = trace.BeginBackendTrace(trace.BackendTrace{Timestamp: startTime, Type: trace.BackendTraceDepth, ModelName: modelConfig.Name, Backend: modelConfig.Backend, Summary: trace.TruncateString(in.GetSrc(), 200)})
4148
}
49+
defer trace.CancelBackendTrace(traceID)
4250

4351
// Stamped here for the same reason as in rerank.go: the caller builds the
4452
// request without a ModelConfig, this function has the one that loaded.
@@ -53,6 +61,7 @@ func Depth(
5361
}
5462

5563
trace.RecordBackendTrace(trace.BackendTrace{
64+
ID: traceID,
5665
Timestamp: startTime,
5766
Duration: time.Since(startTime),
5867
Type: trace.BackendTraceDepth,

core/backend/detection.go

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -33,11 +33,19 @@ func Detection(
3333
return nil, fmt.Errorf("could not load detection model")
3434
}
3535

36+
release, err := AcquireGlobalBackendSlot()
37+
if err != nil {
38+
return nil, err
39+
}
40+
defer release()
3641
var startTime time.Time
42+
var traceID string
3743
if appConfig.EnableTracing {
3844
trace.InitBackendTracingIfEnabled(appConfig.TracingMaxItems, appConfig.TracingMaxBodyBytes)
3945
startTime = time.Now()
46+
traceID = trace.BeginBackendTrace(trace.BackendTrace{Timestamp: startTime, Type: trace.BackendTraceDetection, ModelName: modelConfig.Name, Backend: modelConfig.Backend, Summary: trace.TruncateString(sourceFile, 200)})
4047
}
48+
defer trace.CancelBackendTrace(traceID)
4149

4250
res, err := detectionModel.Detect(ctx, &proto.DetectOptions{
4351
ModelIdentity: modelConfig.Model,
@@ -55,6 +63,7 @@ func Detection(
5563
}
5664

5765
trace.RecordBackendTrace(trace.BackendTrace{
66+
ID: traceID,
5867
Timestamp: startTime,
5968
Duration: time.Since(startTime),
6069
Type: trace.BackendTraceDetection,

core/backend/detokenize.go

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -23,11 +23,19 @@ func ModelDetokenize(tokens []int32, loader *model.ModelLoader, modelConfig conf
2323
return schema.DetokenizeResponse{}, err
2424
}
2525

26+
release, err := AcquireGlobalBackendSlot()
27+
if err != nil {
28+
return schema.DetokenizeResponse{}, err
29+
}
30+
defer release()
2631
var startTime time.Time
32+
var traceID string
2733
if appConfig.EnableTracing {
2834
trace.InitBackendTracingIfEnabled(appConfig.TracingMaxItems, appConfig.TracingMaxBodyBytes)
2935
startTime = time.Now()
36+
traceID = trace.BeginBackendTrace(trace.BackendTrace{Timestamp: startTime, Type: trace.BackendTraceTokenize, ModelName: modelConfig.Name, Backend: modelConfig.Backend, Summary: "detokenize"})
3037
}
38+
defer trace.CancelBackendTrace(traceID)
3139

3240
resp, err := inferenceModel.Detokenize(appConfig.Context, &pb.DetokenizeRequest{Tokens: tokens})
3341

@@ -43,6 +51,7 @@ func ModelDetokenize(tokens []int32, loader *model.ModelLoader, modelConfig conf
4351
}
4452

4553
trace.RecordBackendTrace(trace.BackendTrace{
54+
ID: traceID,
4655
Timestamp: startTime,
4756
Duration: time.Since(startTime),
4857
Type: trace.BackendTraceTokenize,

core/backend/embeddings.go

Lines changed: 17 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -109,9 +109,15 @@ func ModelEmbedding(ctx context.Context, s string, tokens []int, loader *model.M
109109
traceData["input_tokens_count"] = len(tokens)
110110
}
111111

112-
startTime := time.Now()
112+
summary := trace.TruncateString(s, 200)
113+
if summary == "" {
114+
summary = fmt.Sprintf("tokens[%d]", len(tokens))
115+
}
113116
originalFn := wrappedFn
114117
wrappedFn = func() ([]float32, error) {
118+
startTime := time.Now()
119+
traceID := trace.BeginBackendTrace(trace.BackendTrace{Timestamp: startTime, Type: trace.BackendTraceEmbedding, ModelName: modelConfig.Name, Backend: modelConfig.Backend, Summary: summary})
120+
defer trace.CancelBackendTrace(traceID)
115121
result, err := originalFn()
116122
duration := time.Since(startTime)
117123

@@ -122,12 +128,8 @@ func ModelEmbedding(ctx context.Context, s string, tokens []int, loader *model.M
122128
errStr = err.Error()
123129
}
124130

125-
summary := trace.TruncateString(s, 200)
126-
if summary == "" {
127-
summary = fmt.Sprintf("tokens[%d]", len(tokens))
128-
}
129-
130131
trace.RecordBackendTrace(trace.BackendTrace{
132+
ID: traceID,
131133
Timestamp: startTime,
132134
Duration: duration,
133135
Type: trace.BackendTraceEmbedding,
@@ -141,6 +143,15 @@ func ModelEmbedding(ctx context.Context, s string, tokens []int, loader *model.M
141143
return result, err
142144
}
143145
}
146+
originalFn := wrappedFn
147+
wrappedFn = func() ([]float32, error) {
148+
release, err := AcquireGlobalBackendSlot()
149+
if err != nil {
150+
return nil, err
151+
}
152+
defer release()
153+
return originalFn()
154+
}
144155

145156
return wrappedFn, nil
146157
}

core/backend/face_analyze.go

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -30,11 +30,19 @@ func FaceAnalyze(
3030
return nil, fmt.Errorf("could not load face recognition model")
3131
}
3232

33+
release, err := AcquireGlobalBackendSlot()
34+
if err != nil {
35+
return nil, err
36+
}
37+
defer release()
3338
var startTime time.Time
39+
var traceID string
3440
if appConfig.EnableTracing {
3541
trace.InitBackendTracingIfEnabled(appConfig.TracingMaxItems, appConfig.TracingMaxBodyBytes)
3642
startTime = time.Now()
43+
traceID = trace.BeginBackendTrace(trace.BackendTrace{Timestamp: startTime, Type: trace.BackendTraceFaceAnalyze, ModelName: modelConfig.Name, Backend: modelConfig.Backend, Summary: "face analysis"})
3744
}
45+
defer trace.CancelBackendTrace(traceID)
3846

3947
res, err := faceModel.FaceAnalyze(ctx, &proto.FaceAnalyzeRequest{
4048
ModelIdentity: modelConfig.Model,
@@ -49,6 +57,7 @@ func FaceAnalyze(
4957
errStr = err.Error()
5058
}
5159
trace.RecordBackendTrace(trace.BackendTrace{
60+
ID: traceID,
5261
Timestamp: startTime,
5362
Duration: time.Since(startTime),
5463
Type: trace.BackendTraceFaceAnalyze,

core/backend/face_embed.go

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -32,6 +32,11 @@ func FaceEmbed(
3232

3333
predictOpts := gRPCPredictOpts(modelConfig, loader.ModelPath)
3434
predictOpts.Images = []string{imgBase64}
35+
release, err := AcquireGlobalBackendSlot()
36+
if err != nil {
37+
return nil, err
38+
}
39+
defer release()
3540

3641
res, err := faceModel.Embeddings(ctx, predictOpts)
3742
if err != nil {

core/backend/face_verify.go

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -30,11 +30,19 @@ func FaceVerify(
3030
return nil, fmt.Errorf("could not load face recognition model")
3131
}
3232

33+
release, err := AcquireGlobalBackendSlot()
34+
if err != nil {
35+
return nil, err
36+
}
37+
defer release()
3338
var startTime time.Time
39+
var traceID string
3440
if appConfig.EnableTracing {
3541
trace.InitBackendTracingIfEnabled(appConfig.TracingMaxItems, appConfig.TracingMaxBodyBytes)
3642
startTime = time.Now()
43+
traceID = trace.BeginBackendTrace(trace.BackendTrace{Timestamp: startTime, Type: trace.BackendTraceFaceVerify, ModelName: modelConfig.Name, Backend: modelConfig.Backend, Summary: "face verification"})
3744
}
45+
defer trace.CancelBackendTrace(traceID)
3846

3947
res, err := faceModel.FaceVerify(ctx, &proto.FaceVerifyRequest{
4048
ModelIdentity: modelConfig.Model,
@@ -50,6 +58,7 @@ func FaceVerify(
5058
errStr = err.Error()
5159
}
5260
trace.RecordBackendTrace(trace.BackendTrace{
61+
ID: traceID,
5362
Timestamp: startTime,
5463
Duration: time.Since(startTime),
5564
Type: trace.BackendTraceFaceVerify,

core/backend/global_admission.go

Lines changed: 72 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,72 @@
1+
// SPDX-License-Identifier: MIT
2+
3+
package backend
4+
5+
import (
6+
"fmt"
7+
"sync"
8+
"time"
9+
10+
"github.com/mudler/LocalAI/core/config"
11+
)
12+
13+
// BackendAdmissionError reports that the process-wide backend execution
14+
// ceiling is full. HTTP callers map it to 503; internal callers receive the
15+
// same typed error instead of silently queueing and growing in-flight state.
16+
type BackendAdmissionError struct {
17+
Limit int
18+
RetryAfter time.Duration
19+
}
20+
21+
func (e *BackendAdmissionError) Error() string {
22+
return fmt.Sprintf("backend inference capacity reached (max_concurrent=%d); retry after %s", e.Limit, e.RetryAfter)
23+
}
24+
25+
var backendAdmission = struct {
26+
sync.RWMutex
27+
limit int
28+
slots chan struct{}
29+
}{}
30+
31+
// ConfigureGlobalBackendAdmission sets the process-wide ceiling. It is called
32+
// during application construction, before backend work can begin.
33+
func ConfigureGlobalBackendAdmission(limit int) {
34+
if limit <= 0 {
35+
limit = config.DefaultMaxConcurrentBackendRequests
36+
}
37+
backendAdmission.Lock()
38+
backendAdmission.limit = limit
39+
backendAdmission.slots = make(chan struct{}, limit)
40+
backendAdmission.Unlock()
41+
}
42+
43+
// AcquireGlobalBackendSlot admits one backend operation without queueing.
44+
// Callers must invoke release on every completion path.
45+
func AcquireGlobalBackendSlot() (release func(), err error) {
46+
backendAdmission.RLock()
47+
limit, slots := backendAdmission.limit, backendAdmission.slots
48+
backendAdmission.RUnlock()
49+
if slots == nil {
50+
backendAdmission.Lock()
51+
if backendAdmission.slots == nil {
52+
backendAdmission.limit = config.DefaultMaxConcurrentBackendRequests
53+
backendAdmission.slots = make(chan struct{}, backendAdmission.limit)
54+
}
55+
limit, slots = backendAdmission.limit, backendAdmission.slots
56+
backendAdmission.Unlock()
57+
}
58+
select {
59+
case slots <- struct{}{}:
60+
var once sync.Once
61+
return func() { once.Do(func() { <-slots }) }, nil
62+
default:
63+
return nil, &BackendAdmissionError{Limit: limit, RetryAfter: time.Second}
64+
}
65+
}
66+
67+
// GlobalBackendInFlight is the current number of admitted backend operations.
68+
func GlobalBackendInFlight() int {
69+
backendAdmission.RLock()
70+
defer backendAdmission.RUnlock()
71+
return len(backendAdmission.slots)
72+
}

0 commit comments

Comments
 (0)