Files
ThothII/tools/tht/internal/service/service.go
T

199 lines
5.8 KiB
Go

// Package service contains the shared non-mutating health and lifecycle checks for Compose services.
package service
import (
"context"
"encoding/json"
"errors"
"fmt"
"io"
"strings"
"time"
"github.com/aritmolab/thothii/tools/tht/internal/compose"
"github.com/aritmolab/thothii/tools/tht/internal/config"
)
const healthTimeout = 5 * time.Minute
var healthPollInterval = time.Second
// HealthFailure identifies the last non-ready service after a bounded health wait.
type HealthFailure struct {
Service string
State string
}
func (e HealthFailure) Error() string {
return fmt.Sprintf("service health timed out waiting for %s (last state: %s)", e.Service, e.State)
}
type status struct {
Service string `json:"Service"`
State string `json:"State"`
Health string `json:"Health"`
ExitCode json.RawMessage `json:"ExitCode"`
}
// Start optionally builds the configured project, starts it in the background, then waits for
// every required service to become healthy. It does not inspect or modify host state directly.
func Start(ctx context.Context, installation config.Installation, runner compose.Runner, build bool) error {
if runner == nil {
return errors.New("start requires a Docker command runner")
}
if build {
if err := runCompose(ctx, installation, runner, "build"); err != nil {
return fmt.Errorf("image build: %w", err)
}
}
if err := runCompose(ctx, installation, runner, "up", "--detach", "--remove-orphans"); err != nil {
return fmt.Errorf("stack start: %w", err)
}
return WaitForHealthy(ctx, installation, runner)
}
// WaitForHealthy waits for the complete ThothII Compose service set.
func WaitForHealthy(ctx context.Context, installation config.Installation, runner compose.Runner) error {
healthContext, cancel := context.WithTimeout(ctx, healthTimeout)
defer cancel()
lastService, lastState := "core", "unknown"
for {
result, err := runner.Run(healthContext, installation.ComposeArgs("ps", "--all", "--format", "json"), nil)
if err == nil {
statuses, parseErr := parseStatuses(result.Stdout)
if parseErr == nil {
if service, state, ready := healthy(statuses); ready {
return nil
} else {
lastService, lastState = service, state
}
} else {
lastState = "Compose returned invalid service status"
}
} else if result.ExitCode != 0 {
lastState = fmt.Sprintf("Compose exited with status %d", result.ExitCode)
} else {
lastState = "Compose status command failed"
}
select {
case <-healthContext.Done():
return HealthFailure{Service: lastService, State: lastState}
case <-time.After(healthPollInterval):
}
}
}
// CoreRunning reports whether Compose currently identifies core as a running container.
func CoreRunning(value string) (bool, error) {
statuses, err := parseStatuses(value)
if err != nil {
return false, err
}
for _, item := range statuses {
if item.Service == "core" {
return strings.EqualFold(item.State, "running"), nil
}
}
return false, nil
}
// CoreHealthy reports whether core is both running and certified healthy by Compose.
func CoreHealthy(value string) (bool, error) {
statuses, err := parseStatuses(value)
if err != nil {
return false, err
}
for _, item := range statuses {
if item.Service == "core" {
return strings.EqualFold(item.State, "running") && strings.EqualFold(item.Health, "healthy"), nil
}
}
return false, nil
}
// Healthy verifies the full expected Compose service set, including the one-shot model initializer.
func Healthy(value string) error {
statuses, err := parseStatuses(value)
if err != nil {
return err
}
service, state, ready := healthy(statuses)
if ready {
return nil
}
return fmt.Errorf("%s is not healthy (%s)", service, state)
}
func runCompose(ctx context.Context, installation config.Installation, runner compose.Runner, command ...string) error {
result, err := runner.Run(ctx, installation.ComposeArgs(command...), nil)
if err == nil {
return nil
}
if result.ExitCode != 0 {
return fmt.Errorf("Docker exited with status %d", result.ExitCode)
}
return err
}
func parseStatuses(value string) ([]status, error) {
var statuses []status
if err := json.Unmarshal([]byte(value), &statuses); err == nil && len(statuses) > 0 {
return statuses, nil
}
decoder := json.NewDecoder(strings.NewReader(value))
for {
var item status
err := decoder.Decode(&item)
if errors.Is(err, io.EOF) {
break
}
if err != nil {
return nil, errors.New("Compose returned invalid service status")
}
statuses = append(statuses, item)
}
if len(statuses) == 0 {
return nil, errors.New("Compose returned invalid service status")
}
return statuses, nil
}
func healthy(statuses []status) (string, string, bool) {
wanted := map[string]bool{"core": false, "frontend": false, "qdrant": false, "embedding": false, "embedding-model-init": false}
for _, item := range statuses {
if _, required := wanted[item.Service]; !required {
continue
}
if item.Service == "embedding-model-init" {
if strings.EqualFold(item.State, "exited") && exitCodeZero(item.ExitCode) {
wanted[item.Service] = true
continue
}
return item.Service, item.State, false
}
if strings.EqualFold(item.State, "running") && strings.EqualFold(item.Health, "healthy") {
wanted[item.Service] = true
continue
}
return item.Service, strings.TrimSpace(item.State + "/" + item.Health), false
}
for _, name := range []string{"core", "frontend", "qdrant", "embedding", "embedding-model-init"} {
if !wanted[name] {
return name, "not reported by Docker Compose", false
}
}
return "", "", true
}
func exitCodeZero(value json.RawMessage) bool {
if len(value) == 0 || string(value) == "null" {
return false
}
var number int
if json.Unmarshal(value, &number) == nil {
return number == 0
}
var text string
return json.Unmarshal(value, &text) == nil && strings.TrimSpace(text) == "0"
}