Repository navigation
Expand file tree
/
Copy pathmain.go
More file actions
470 lines (434 loc) · 16.7 KB
/
Copy pathmain.go
File metadata and controls
470 lines (434 loc) · 16.7 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
// qoder2api: OpenAI-compatible API bridge for QoderWork.
package main
import (
"context"
"encoding/json"
"log"
"net/http"
"os"
"os/signal"
"strconv"
"syscall"
"qoder2api/admin"
"qoder2api/auth"
"qoder2api/bridge"
"qoder2api/logs"
"qoder2api/models"
"qoder2api/stats"
"qoder2api/store"
"qoder2api/web"
"sync"
"time"
)
// Version is the build version, overridden at build time via
// -ldflags "-X main.Version=v1.2.3". It is only used for display (login page).
var Version = "dev"
// bridgeProvider manages a single OpenAiBridge for the current PAT.
// When the PAT changes (via admin UI), the bridge is recreated on next access.
type bridgeProvider struct {
mu sync.Mutex
bridge *bridge.OpenAiBridge
pat string
store *store.Store
// onPatChange is fired after a different non-empty PAT was replaced or the
// PAT was cleared. It runs without the provider lock held: the callback
// takes the recorder lock, and nesting the two risks a lock-order stall.
onPatChange func()
}
func newBridgeProvider(st *store.Store) *bridgeProvider {
return &bridgeProvider{store: st}
}
// setOnPatChange installs the callback fired when the configured PAT changes.
// Wiring happens before the provider is used by any goroutine.
func (p *bridgeProvider) setOnPatChange(fn func()) {
p.mu.Lock()
p.onPatChange = fn
p.mu.Unlock()
}
// resolveBridge validates the API key and returns the bridge for the current
// PAT plus the key's metadata for request logging. Returns nil if the key is
// invalid or no PAT is configured.
func (p *bridgeProvider) resolveBridge(apiKey string) *bridge.ResolvedBridge {
key, ok := p.store.LookupKey(apiKey)
if !ok {
return nil
}
b := p.currentBridge()
if b == nil {
return nil
}
return &bridge.ResolvedBridge{Bridge: b, KeyID: key.ID, KeyNote: key.Note}
}
// currentBridge returns the bridge for the current PAT, creating or recreating
// the bridge when the PAT changes. A PAT change (including clearing it) also
// fires onPatChange after the provider lock is released, so the cached account
// snapshot of the previous PAT can be dropped.
func (p *bridgeProvider) currentBridge() *bridge.OpenAiBridge {
pat := p.store.GetPAT()
p.mu.Lock()
changed := p.pat != "" && p.pat != pat
var b *bridge.OpenAiBridge
if pat == "" {
if changed {
// The old PAT/session must not linger once it is cleared.
p.bridge = nil
p.pat = ""
}
} else if p.bridge == nil || p.pat != pat {
realPAT, region := auth.Resolve(pat)
p.bridge = bridge.NewOpenAiBridge(realPAT, region)
p.pat = pat
b = p.bridge
} else {
b = p.bridge
}
hook := p.onPatChange
p.mu.Unlock()
if changed && hook != nil {
hook()
}
return b
}
// accountRefreshCadence keeps the authoritative allowance reasonably fresh
// without coupling the admin panel's 15-second poll to an upstream request.
const accountRefreshCadence = time.Minute
// patWriteMu serializes the recorder writes of the account loop against the
// PAT-change clear. Without it, an in-flight refresh could pass its PAT
// re-check, be preempted by an admin PAT change, and then write the previous
// account's snapshot after the hook had already cleared it. The lock is never
// held across the upstream fetch.
var patWriteMu sync.Mutex
// applyAccountWrite persists one account snapshot under patWriteMu, re-checking
// the PAT inside the lock so a snapshot belonging to a replaced or cleared PAT
// can never be written. It returns false when the snapshot was dropped.
func applyAccountWrite(rec *stats.Recorder, provider *bridgeProvider, pat string, nextResetMs int64, account *stats.Account) bool {
patWriteMu.Lock()
defer patWriteMu.Unlock()
if provider.store.GetPAT() != pat {
return false // the PAT changed while the fetch was in flight
}
rec.SetBillingCycle(nextResetMs)
rec.SetAccount(account)
return true
}
// clearAccountOnPatChange drops the cached account snapshot when the configured
// PAT changed. It shares patWriteMu with applyAccountWrite so the clear and the
// writes are mutually exclusive.
func clearAccountOnPatChange(rec *stats.Recorder) {
patWriteMu.Lock()
defer patWriteMu.Unlock()
rec.ClearAccount()
}
// healthHandler serves an unauthenticated liveness/readiness probe. It is
// deliberately cheap (no upstream calls, no session bootstrap): the body
// carries readiness details (PAT configured, effective stream timeouts)
// while the status stays 200 as long as the process is serving.
func healthHandler(st *store.Store) http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodGet && r.Method != http.MethodHead {
w.Header().Set("Allow", "GET, HEAD")
http.Error(w, "method not allowed", http.StatusMethodNotAllowed)
return
}
timeouts := auth.CurrentStreamTimeouts()
w.Header().Set("Content-Type", "application/json")
json.NewEncoder(w).Encode(map[string]interface{}{
"status": "ok",
"has_pat": st.GetPAT() != "",
"chat_timeout_seconds": int(timeouts.Header.Seconds()),
"idle_timeout_seconds": int(timeouts.Idle.Seconds()),
})
}
}
func main() {
// Initialize store
dataPath := os.Getenv("QODER_DATA_PATH")
if dataPath == "" {
dataPath = "data.json"
}
st, err := store.New(dataPath)
if err != nil {
log.Fatalf("[store] failed to init: %v", err)
}
// Host/port: env vars override config file; config file overrides defaults
host := os.Getenv("QODER_HOST")
if host == "" {
host = st.GetHost()
}
port := st.GetPort()
if envPort := os.Getenv("QODER_PORT"); envPort != "" {
if p, err := strconv.Atoi(envPort); err == nil {
port = p
} else {
log.Printf("[bridge] WARN invalid QODER_PORT=%q; using %d", envPort, port)
}
}
// Chat stream timeouts: env vars override the persisted config, the
// config file overrides the built-in defaults (same precedence as
// host/port). Applied to new requests immediately.
chatTimeout := st.GetChatTimeoutSeconds()
if v := os.Getenv("QODER_CHAT_TIMEOUT_SECONDS"); v != "" {
if n, err := strconv.Atoi(v); err == nil && n >= 1 && n <= store.MaxChatTimeoutSeconds {
chatTimeout = n
} else {
log.Printf("[bridge] WARN invalid QODER_CHAT_TIMEOUT_SECONDS=%q; using %d", v, chatTimeout)
}
}
idleTimeout := st.GetIdleTimeoutSeconds()
if v := os.Getenv("QODER_IDLE_TIMEOUT_SECONDS"); v != "" {
if n, err := strconv.Atoi(v); err == nil && n >= 1 && n <= store.MaxIdleTimeoutSeconds {
idleTimeout = n
} else {
log.Printf("[bridge] WARN invalid QODER_IDLE_TIMEOUT_SECONDS=%q; using %d", v, idleTimeout)
}
}
auth.SetStreamTimeouts(time.Duration(chatTimeout)*time.Second, time.Duration(idleTimeout)*time.Second)
// Request-log limits: env vars override the persisted config, which
// overrides the built-in defaults (same precedence as the timeouts above).
logRetention := st.GetLogRetentionDays()
if v := os.Getenv("QODER_LOG_RETENTION_DAYS"); v != "" {
if n, err := strconv.Atoi(v); err == nil && n >= 1 && n <= store.MaxLogRetentionDays {
logRetention = n
} else {
log.Printf("[logs] WARN invalid QODER_LOG_RETENTION_DAYS=%q; using %d", v, logRetention)
}
}
logMaxEntries := st.GetLogMaxEntries()
if v := os.Getenv("QODER_LOG_MAX_ENTRIES"); v != "" {
if n, err := strconv.Atoi(v); err == nil && n >= store.MinLogMaxEntries && n <= store.MaxLogEntriesHardCap {
logMaxEntries = n
} else {
log.Printf("[logs] WARN invalid QODER_LOG_MAX_ENTRIES=%q; using %d", v, logMaxEntries)
}
}
provider := newBridgeProvider(st)
// Request statistics recorder with periodic persistence (30s cadence).
// The loop also drains a stop channel so graceful shutdown can trigger a
// final flush, avoiding loss of the last <=30s of stats.
rec := stats.NewRecorder(st)
// Request-log recorder: same persistence model as stats (in-memory on the
// request path, flushed out of band). Entries are metadata only.
logRec := logs.NewRecorder(st, logRetention, logMaxEntries)
// A PAT change invalidates the previous account's plan and allowance; drop
// the cached snapshot so the panel cannot keep showing it. The clear shares
// patWriteMu with the account loop's writes.
provider.setOnPatChange(func() { clearAccountOnPatChange(rec) })
// The snapshot is persisted, so a restart with a different or cleared PAT
// would otherwise keep showing the previous account (the hook only fires on
// an in-process change). Drop it once here; the account loop repopulates the
// panel from the first successful refresh.
rec.ClearAccount()
statsStop := make(chan struct{})
statsDone := make(chan struct{})
go func() {
defer close(statsDone)
ticker := time.NewTicker(30 * time.Second)
defer ticker.Stop()
flushAll := func(final bool) {
if err := rec.Flush(); err != nil {
log.Printf("[stats] WARN %sflush failed: %v", finalStr(final), err)
}
if err := logRec.Flush(); err != nil {
log.Printf("[logs] WARN %sflush failed: %v", finalStr(final), err)
}
}
for {
select {
case <-ticker.C:
flushAll(false)
case <-statsStop:
flushAll(true)
return
}
}
}()
// Model fetcher for admin UI
modelFetcher := func(ctx context.Context) []string {
b := provider.currentBridge()
if b == nil {
return models.DefaultCatalog().Keys()
}
catalog := b.GetCatalog(ctx)
return catalog.Keys()
}
// Subscription account loop: refreshes identity metadata from /user/status
// and authoritative allowance totals from OpenAPI /api/v2/quota/usage out
// of band. The admin panel's 15-second poll remains memory-only.
//
// The loop has its own stop/done pair so shutdown can join it before the
// final stats flush: its last SetAccount/SetBillingCycle would otherwise
// land after that flush and be lost on exit.
accountStop := make(chan struct{})
accountDone := make(chan struct{})
applyAccount := func(ctx context.Context) {
// Capture the PAT this refresh belongs to: if it changes while the
// fetch is in flight, the snapshot must be dropped rather than written
// after the PAT-change hook cleared the previous account.
pat := provider.store.GetPAT()
b := provider.currentBridge()
if b == nil {
return // no PAT configured yet
}
st := b.EnsureAccountStatus(ctx)
if st == nil {
return // upstream unreachable: keep the last known state
}
account := &stats.Account{
Plan: st.Plan,
Tag: st.UserTag,
OrgName: st.OrgName,
IsQuotaExceeded: st.IsQuotaExceeded,
TotalUsagePercentage: st.TotalUsagePercentage,
}
if st.UserQuota != nil {
account.UserQuota = &stats.Quota{
Total: st.UserQuota.Total, Used: st.UserQuota.Used,
Remaining: st.UserQuota.Remaining, Percentage: st.UserQuota.Percentage,
Unit: st.UserQuota.Unit, DetailURL: st.UserQuota.DetailURL,
}
}
if st.AddOnQuota != nil {
account.AddOnQuota = &stats.Quota{
Total: st.AddOnQuota.Total, Used: st.AddOnQuota.Used,
Remaining: st.AddOnQuota.Remaining, Percentage: st.AddOnQuota.Percentage,
Unit: st.AddOnQuota.Unit, DetailURL: st.AddOnQuota.DetailURL,
}
}
if st.OrgResourcePackage != nil {
account.OrgResourcePackage = &stats.OrgResourcePackage{
Used: st.OrgResourcePackage.Used, Cap: st.OrgResourcePackage.Cap,
Remaining: st.OrgResourcePackage.Remaining, Percentage: st.OrgResourcePackage.Percentage,
Available: st.OrgResourcePackage.Available, Unit: st.OrgResourcePackage.Unit,
}
}
// The write section runs under patWriteMu and re-checks the PAT inside it:
// the PAT-change hook clears the recorder under the same lock, so a
// snapshot fetched for the previous PAT can no longer land after that
// clear. The fetch above stays outside the lock (it can take 30s).
applyAccountWrite(rec, provider, pat, st.NextResetAtMs, account)
}
go func() {
defer close(accountDone)
// Fetch once immediately so the panel is populated from the start
// rather than after the first tick.
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
applyAccount(ctx)
cancel()
ticker := time.NewTicker(accountRefreshCadence)
defer ticker.Stop()
for {
select {
case <-ticker.C:
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
applyAccount(ctx)
cancel()
case <-accountStop:
return
}
}
}()
// Initialize admin
adminInst := admin.New(st, modelFetcher, rec, logRec)
// WebUI: pre-render all pages at startup (multi-page Alpine.js frontend,
// no build step, no external dependencies).
webHandler, err := web.New(Version)
if err != nil {
log.Fatalf("[web] failed to init: %v", err)
}
mux := http.NewServeMux()
mux.HandleFunc("/v1/chat/completions", bridge.MakeChatHandler(provider.resolveBridge, rec, logRec))
mux.HandleFunc("/v1/models", bridge.MakeModelsHandler(provider.resolveBridge))
mux.HandleFunc("/health", healthHandler(st))
// Root redirect to admin UI
mux.HandleFunc("/", func(w http.ResponseWriter, r *http.Request) {
if r.URL.Path == "/" {
http.Redirect(w, r, "/admin", http.StatusFound)
return
}
http.NotFound(w, r)
})
// Register admin API routes (/admin/api/*) and WebUI routes
// (/admin, /admin/<page>, /admin/assets/*)
adminInst.RegisterRoutes(mux)
webHandler.Mount(mux)
addr := host + ":" + strconv.Itoa(port)
log.Printf("[bridge] listening http://%s/v1/chat/completions", addr)
log.Printf("[admin] http://%s/admin", addr)
// Timeouts: ReadHeaderTimeout/IdleTimeout harden against slowloris-style
// resource exhaustion. WriteTimeout is intentionally unset so that
// long-lived SSE streaming responses are not prematurely cancelled.
srv := &http.Server{
Addr: addr,
Handler: mux,
ReadHeaderTimeout: 10 * time.Second,
IdleTimeout: 120 * time.Second,
}
// Graceful shutdown on SIGINT/SIGTERM: stop the account refresher and join
// it, trigger the final stats flush, and drain in-flight requests with a
// bounded deadline.
//
// stopOnce is shared by every caller of shutdown/stopBackgroundLoops so a
// second signal (or a future forced-exit path) cannot panic on a double
// close of the stop channels.
stopOnce := &sync.Once{}
shutdownErr := make(chan error, 1)
go func() {
sigCh := make(chan os.Signal, 1)
signal.Notify(sigCh, syscall.SIGINT, syscall.SIGTERM)
sig := <-sigCh
// Unregister before tearing down: restoring the default disposition
// means a second signal terminates the process immediately instead of
// being queued into a handler that has already run. The stop channels
// themselves are guarded by sync.Once, so a duplicate teardown is safe.
signal.Stop(sigCh)
log.Printf("[server] received %s, shutting down...", sig)
shutdown(shutdownErr, srv, stopOnce, accountStop, statsStop, accountDone, statsDone)
}()
if err := srv.ListenAndServe(); err != nil && err != http.ErrServerClosed {
log.Fatalf("[server] failed to listen on %s: %v\n[server] If port %d is in use, change it in the admin UI or data.json", addr, err, port)
}
if err := <-shutdownErr; err != nil {
log.Printf("[server] shutdown error: %v", err)
}
// Wait for the stats goroutine to finish its final flush before exiting.
<-statsDone
log.Printf("[server] stopped")
}
// finalStr labels a flush log line ("final " on shutdown) so periodic and
// shutdown flushes stay distinguishable.
func finalStr(final bool) string {
if final {
return "final "
}
return ""
}
// shutdown stops the background loops and drains in-flight requests. The drain
// runs concurrently with the join so a slow account refresh cannot delay it, but
// it is still awaited before this returns: cancelling the drain's context early
// would abort every in-flight streaming response on SIGTERM, which is what a
// previous version of this function did. The stop channels are closed through
// stopOnce, so a repeated shutdown call cannot panic.
func shutdown(shutdownErr chan<- error, srv *http.Server, stopOnce *sync.Once, accountStop, statsStop chan<- struct{}, accountDone, statsDone <-chan struct{}) {
ctx, cancel := context.WithTimeout(context.Background(), 15*time.Second)
drain := make(chan error, 1)
go func() { drain <- srv.Shutdown(ctx) }()
stopBackgroundLoops(stopOnce, accountStop, statsStop, accountDone, statsDone)
shutdownErr <- <-drain
cancel()
}
// stopBackgroundLoops joins the account refresher before triggering the stats
// loop's final flush: the account loop may be mid-fetch for up to 30s, and a
// SetAccount/SetBillingCycle landing after that flush would be lost on exit.
//
// once guards every stop-channel close, so a second caller (a future
// "second signal forces exit" path, or a duplicate call) blocks until the
// first stop completes and then returns instead of panicking on a double
// close. Callers supply the Once so independent instances stay independent.
func stopBackgroundLoops(once *sync.Once, accountStop, statsStop chan<- struct{}, accountDone, statsDone <-chan struct{}) {
once.Do(func() {
close(accountStop)
<-accountDone
close(statsStop)
<-statsDone
})
}