diff --git a/deploy/README.md b/deploy/README.md index 200b6803..cc2a9029 100644 --- a/deploy/README.md +++ b/deploy/README.md @@ -9,7 +9,7 @@ ## Installation -`install.sh` downloads its release's `compose.yaml`, checks it against `compose-sha256sums.txt`, writes `.env`, and starts Compose. Core applies database migrations when it starts. The host needs Linux amd64 and Docker Compose 2.26 or newer. [Configuration](../docs/configuration.md) owns the installation layout and settings. +`install.sh` downloads its release's `compose.yaml`, checks it against `compose-sha256sums.txt`, writes `.env`, and starts Compose. Installation commands stream their progress and report elapsed time; [startup diagnostics](../docs/getting-started/operations.md#startup-diagnostics) describes service logs. Core applies database migrations when it starts. The host needs Linux amd64 and Docker Compose 2.26 or newer. [Configuration](../docs/configuration.md) owns the installation layout and settings. `oac` is a Go command (`services/core/cmd/oac`) in the Core image and the ingress image. The host copy implements `apply`, `core-key` and `rotate-core-key`; `core-key --show` runs `oac-web core-key` in the Web container. Start, stop, logs and removal are `docker compose`. `apply` runs `oac-core check-config` before recreating services. The ingress image runs data initialization as `oac init`, verifies and copies its bundled node metadata without network access, and contains no Python. No service receives a Docker socket. diff --git a/deploy/compose/compose.yaml b/deploy/compose/compose.yaml index 21ecceea..7b613632 100644 --- a/deploy/compose/compose.yaml +++ b/deploy/compose/compose.yaml @@ -11,6 +11,9 @@ services: command: [/usr/local/bin/oac, init] environment: OAC_REVISION: __OAC_REVISION__ + OAC_LOG_LEVEL: ${OAC_LOG_LEVEL:-} + OAC_LOG_FORMAT: ${OAC_LOG_FORMAT:-} + OAC_LOG_ADD_SOURCE: ${OAC_LOG_ADD_SOURCE:-} volumes: - type: bind source: ${OAC_DATA_DIR:-./data} diff --git a/deploy/install.sh b/deploy/install.sh index f862fb28..013d1289 100755 --- a/deploy/install.sh +++ b/deploy/install.sh @@ -73,7 +73,6 @@ if [[ "$version" != latest ]]; then asset_base="https://github.com/${repository}/releases/download/${version}" fi -log="$(mktemp)" cleanup() { if [[ "$kept" != 1 && -d "$install_dir" ]]; then ( @@ -82,18 +81,23 @@ cleanup() { docker compose down --remove-orphans # Containers own data/; remove it from a container as well. if [[ -d data ]]; then docker compose run --rm --no-deps --volume "$install_dir/data:/data" --entrypoint find database /data -mindepth 1 -delete; fi - ) >/dev/null 2>&1 || true + ) || true rm -rf "$install_dir" fi - rm -f "$log" } trap cleanup EXIT -# step DESCRIPTION COMMAND... prints the command's output only when it fails. +# step DESCRIPTION COMMAND... streams progress and reports elapsed time. step() { - printf '%s... ' "$1" + local description="$1" step_started="$SECONDS" + printf '%s...\n' "$description" shift - if "$@" >"$log" 2>&1; then echo done; else echo failed; cat "$log" >&2; return 1; fi + if "$@"; then + printf '%s completed (%ss).\n' "$description" "$((SECONDS - step_started))" + else + printf '%s failed (%ss).\n' "$description" "$((SECONDS - step_started))" >&2 + return 1 + fi } # The source address of this host's default route, when it is a private one. @@ -131,7 +135,6 @@ step "Starting services" docker compose up -d --wait key="$(./oac core-key --show)" kept=1 trap - EXIT -rm -f "$log" sudo="" if [[ "$EUID" == 0 && -n "${SUDO_USER:-}" ]]; then sudo="sudo "; fi diff --git a/deploy/test_install.py b/deploy/test_install.py index 6393b28f..4889a0e7 100644 --- a/deploy/test_install.py +++ b/deploy/test_install.py @@ -22,7 +22,7 @@ def install(self, root, *args, compose_up=0, route="1.1.1.1 via 10.0.0.1 dev eth printf '%s\\n' "$*" >> {log} if [ "$1" = compose ] && [ "$2" = version ]; then printf 'v2.29.1\\n'; exit 0; fi if [ "$1" = compose ] && [ "$2" = cp ]; then printf '#!/bin/sh\\necho oac_core_fixture\\n' > ./oac; chmod +x ./oac; exit 0; fi - if [ "$1" = compose ] && [ "$2" = up ]; then exit {compose_up}; fi + if [ "$1" = compose ] && [ "$2" = up ]; then printf 'fixture startup progress\\n'; exit {compose_up}; fi exit 0 """)) self.write_executable(bin_dir / "curl", textwrap.dedent("""\ @@ -61,6 +61,8 @@ def test_the_private_address_is_the_default_public_url(self): self.assertIn("OAC_PUBLIC_URL=http://10.0.0.5:8080\n", (root / "oac/.env").read_text()) self.assertIn("Console http://10.0.0.5:8080", completed.stdout) self.assertIn("Core key oac_core_fixture", completed.stdout) + self.assertIn("fixture startup progress", completed.stdout) + self.assertIn("Starting services completed (", completed.stdout) self.assertNotIn("Only this host", completed.stdout) def test_without_a_private_address_only_this_host_reaches_web(self): diff --git a/docs/getting-started/operations.md b/docs/getting-started/operations.md index b692794b..366601d7 100644 --- a/docs/getting-started/operations.md +++ b/docs/getting-started/operations.md @@ -40,11 +40,27 @@ Service health does not show that a harness or a model works. Use Session, Turn, ```sh docker compose -f "$HOME/.oac/core/compose.yaml" ps --all -docker compose -f "$HOME/.oac/core/compose.yaml" logs --tail 200 core +docker compose -f "$HOME/.oac/core/compose.yaml" logs --tail 200 init database core ``` Don't paste `docker compose config`, `docker inspect` or raw logs into public issue reports. +### Startup diagnostics + +At the default `info` level, initialization and Core log `Stage started`, `Stage completed` and `Stage failed` with `component`, `stage` and completion `elapsed_ms`. Initialization stages cover directories, the installation lock, existing data checks, metadata verification/publication, secrets and the installation receipt. Core stages cover configuration, database migrations/connection, services, Runtime setup, the execution worker and the HTTP listener. `Core HTTP listener ready` means its socket is bound; a failure in `running` occurs after startup. A normal signal starts `shutdown`. + +Configuration validation also identifies the setting and an authored reason without repeating its value. Failures report typed `error_kind` facts without arbitrary error text. Filesystem failures include `operation`, `path` and numeric `errno`; database errors can include `sqlstate`. File contents, credentials and database connection strings are excluded. `unclassified` means no supported typed cause was available; use the stage and adjacent service logs to investigate. + +| `error_kind` | Check | +| --- | --- | +| `not_found`, `permission_denied` | The named file, its mount source, ownership and permissions | +| `connection_refused`, `connection_reset` | The dependency's container status and logs | +| `timeout`, `canceled` | Dependency availability, elapsed time and shutdown events | +| `unexpected_eof`, `eof` | The stage and dependency logs; an interrupted stream alone does not identify its cause | +| `database_error` | PostgreSQL logs and the reported `sqlstate` | + +The host installer streams command progress and reports each step's elapsed seconds, including failures. For file preparation details, set `OAC_LOG_LEVEL=debug` and follow [process settings](../configuration.md#process-settings-configjson). Initialization diagnostics are emitted when its container runs; an already completed one-time container does not run again merely because Core restarts. + ## Stop and restart Let active work settle before a planned restart: diff --git a/docs/zh/getting-started/operations.md b/docs/zh/getting-started/operations.md index facb1309..a40ea94a 100644 --- a/docs/zh/getting-started/operations.md +++ b/docs/zh/getting-started/operations.md @@ -1,7 +1,7 @@ --- title: "管理你的安装" source: docs/getting-started/operations.md -source_hash: e60a6e96cd61692b6adc7664874ba0ed2ce489982e7da2fe2347631ee2a94c0a +source_hash: 3158dbf2ec3254f22fe832cd1113b23137eb2f487fb6c04b828d73e40945cf68 --- 安装运维人员负责 Core 主机、存储和可用性。节点主机运行各自的服务;参阅[节点](nodes.md)。设置见[配置参考](../configuration.md)。 @@ -42,11 +42,27 @@ docker compose -f ~/.oac/core/compose.yaml ps ```sh docker compose -f "$HOME/.oac/core/compose.yaml" ps --all -docker compose -f "$HOME/.oac/core/compose.yaml" logs --tail 200 core +docker compose -f "$HOME/.oac/core/compose.yaml" logs --tail 200 init database core ``` 不要将 `docker compose config`、`docker inspect` 或原始日志粘贴到公开问题报告。 +### 启动诊断 {#startup-diagnostics} + +在默认 `info` 级别下,初始化和 Core 会记录 `Stage started`、`Stage completed` 与 `Stage failed`,包含 `component`、`stage` 以及完成时的 `elapsed_ms`。初始化阶段覆盖目录、安装锁、已有数据检查、元数据验证与发布、机密信息以及安装记录。Core 阶段覆盖配置、数据库迁移与连接、服务、Runtime 配置、执行工作线程以及 HTTP 监听器。`Core HTTP listener ready` 表示套接字已绑定;`running` 阶段的失败发生在启动完成后。正常信号会开始 `shutdown`。 + +配置校验还会指出配置项和明确的原因,不重复其值。失败日志通过类型化的 `error_kind` 描述原因,不输出任意错误文本。文件系统失败包含 `operation`、`path` 和数字 `errno`;数据库错误可包含 `sqlstate`。日志不包含文件内容、凭据或数据库连接字符串。`unclassified` 表示没有可识别的类型化原因;结合阶段及相邻服务日志排查。 + +| `error_kind` | 检查内容 | +| --- | --- | +| `not_found`、`permission_denied` | 指定文件、挂载来源、所有者与权限 | +| `connection_refused`、`connection_reset` | 依赖服务的容器状态和日志 | +| `timeout`、`canceled` | 依赖服务可用性、耗时和停止事件 | +| `unexpected_eof`、`eof` | 阶段和依赖服务日志;流被中断本身不能确定原因 | +| `database_error` | PostgreSQL 日志及记录的 `sqlstate` | + +宿主机安装器实时显示命令进度,并记录每一步的耗时秒数,包括失败步骤。要查看文件准备细节,设置 `OAC_LOG_LEVEL=debug` 并遵循[进程设置](../configuration.md#process-settings-configjson)。初始化诊断在其容器运行时输出;已经完成的一次性容器不会仅因 Core 重启而再次运行。 + ## 停止与重启 {#stop-and-restart} 计划重启前,先等待活动工作结束: diff --git a/internal/obs/log/stage.go b/internal/obs/log/stage.go new file mode 100644 index 00000000..6d796147 --- /dev/null +++ b/internal/obs/log/stage.go @@ -0,0 +1,77 @@ +package log + +import ( + "context" + "errors" + "io" + "io/fs" + "net" + "os" + "syscall" + "time" +) + +// StartStage reports progress without logging operation inputs or error text. +// Call the returned function once with the stage's result. +func StartStage(stage string, fields ...any) func(error) { + logger := With(fields...).With("stage", stage) + started := time.Now() + logger.Info("Stage started") + return func(err error) { + elapsed := time.Since(started).Milliseconds() + if err != nil { + logger.Error("Stage failed", append([]any{"elapsed_ms", elapsed}, ErrorFields(err)...)...) + } else { + logger.Info("Stage completed", "elapsed_ms", elapsed) + } + } +} + +// ErrorFields preserves typed diagnostic facts, never arbitrary error messages. +func ErrorFields(err error) []any { + kind := "unclassified" + switch { + case errors.Is(err, context.Canceled): + kind = "canceled" + case errors.Is(err, context.DeadlineExceeded): + kind = "timeout" + case errors.Is(err, fs.ErrNotExist): + kind = "not_found" + case errors.Is(err, fs.ErrPermission): + kind = "permission_denied" + case errors.Is(err, syscall.ECONNREFUSED): + kind = "connection_refused" + case errors.Is(err, syscall.ECONNRESET): + kind = "connection_reset" + case errors.Is(err, io.ErrUnexpectedEOF): + kind = "unexpected_eof" + case errors.Is(err, io.EOF): + kind = "eof" + } + var timeout net.Error + if kind == "unclassified" && errors.As(err, &timeout) && timeout.Timeout() { + kind = "timeout" + } + fields := []any{"error_kind", kind} + var path *os.PathError + if errors.As(err, &path) { + fields = append(fields, "operation", path.Op, "path", path.Path) + } + var errno syscall.Errno + if errors.As(err, &errno) { + fields = append(fields, "errno", int(errno)) + } + var sql interface{ SQLState() string } + if errors.As(err, &sql) { + state := sql.SQLState() + valid := len(state) == 5 + for _, c := range state { + valid = valid && (c >= '0' && c <= '9' || c >= 'A' && c <= 'Z') + } + if valid { + fields[1] = "database_error" + fields = append(fields, "sqlstate", state) + } + } + return fields +} diff --git a/internal/obs/log/stage_test.go b/internal/obs/log/stage_test.go new file mode 100644 index 00000000..92ef0f17 --- /dev/null +++ b/internal/obs/log/stage_test.go @@ -0,0 +1,67 @@ +package log + +import ( + "bytes" + "context" + "encoding/json" + "errors" + "fmt" + "io" + "log/slog" + "os" + "strings" + "syscall" + "testing" +) + +type databaseFailure struct{ state string } + +func (e databaseFailure) Error() string { return "private database credentials" } +func (e databaseFailure) SQLState() string { return e.state } + +func TestStageFailureReportsTypedFactsWithoutErrorText(t *testing.T) { + cases := []struct { + name string + err error + kind string + }{ + {"permission", &os.PathError{Op: "open", Path: "/run/oac/installation.id", Err: syscall.EACCES}, "permission_denied"}, + {"missing", fmt.Errorf("private credentials: %w", os.ErrNotExist), "not_found"}, + {"refused", fmt.Errorf("private credentials: %w", syscall.ECONNREFUSED), "connection_refused"}, + {"truncated", io.ErrUnexpectedEOF, "unexpected_eof"}, + {"timeout", context.DeadlineExceeded, "timeout"}, + {"database", databaseFailure{"40P01"}, "database_error"}, + {"unknown", errors.New("private credentials"), "unclassified"}, + {"invalid SQL state", databaseFailure{"private credentials"}, "unclassified"}, + } + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + var output bytes.Buffer + previous := slog.Default() + slog.SetDefault(buildLogger(Config{Out: &output, Format: "json"})) + defer slog.SetDefault(previous) + finish := StartStage("database_connection", "component", "core") + finish(tc.err) + if strings.Contains(output.String(), "private") { + t.Fatalf("error text leaked: %s", output.String()) + } + lines := strings.Split(strings.TrimSpace(output.String()), "\n") + if len(lines) != 2 { + t.Fatalf("logs = %s", output.String()) + } + var event map[string]any + if err := json.Unmarshal([]byte(lines[1]), &event); err != nil { + t.Fatal(err) + } + if event["msg"] != "Stage failed" || event["stage"] != "database_connection" || event["error_kind"] != tc.kind || event["elapsed_ms"] == nil { + t.Fatalf("event = %v", event) + } + if tc.name == "permission" && (event["path"] != "/run/oac/installation.id" || event["operation"] != "open") { + t.Fatalf("path facts = %v", event) + } + if tc.name == "database" && event["sqlstate"] != "40P01" { + t.Fatalf("SQL state = %v", event) + } + }) + } +} diff --git a/scripts/compose-smoke.py b/scripts/compose-smoke.py index 2a6b27fa..7cb564ef 100644 --- a/scripts/compose-smoke.py +++ b/scripts/compose-smoke.py @@ -114,7 +114,7 @@ def main(): override.write_text(json.dumps({'services': {'init': {'network_mode': 'none'}}})) env = {**os.environ, 'COMPOSE_PROGRESS': 'plain', 'OAC_DATA_DIR': str(data), - 'OAC_HOST': '127.0.0.1', 'OAC_WEB_PORT': '0', + 'OAC_HOST': '127.0.0.1', 'OAC_WEB_PORT': '0', 'OAC_LOG_LEVEL': 'info', 'OAC_LOG_FORMAT': 'json', **{'OAC_IMAGE_' + name.upper(): image for name, image in images.items()}} env.pop('OAC_PUBLIC_URL', None) command = ['docker', 'compose', '--env-file', os.devnull, '-p', project, @@ -193,7 +193,9 @@ def terminate(_signum, _frame): 'Content-Type: text/plain\r\n\r\n').encode() + content + f'\r\n--{boundary}--\r\n'.encode() uploaded = get('/v1/files', body=body, headers={**api, 'Content-Type': 'multipart/form-data; boundary=' + boundary}) assert uploaded['bytes'] == len(content), 'Upload was truncated' - private_logs(key, project_key) + initial_logs = private_logs(key, project_key) + assert '"stage":"metadata_verification"' in initial_logs, 'Initialization progress is missing' + assert '"msg":"Core HTTP listener ready"' in initial_logs, 'Core readiness diagnostic is missing' print('Configuring a reachable URL and recreating containers with the same data directory', flush=True) # Retain the assigned port across recreation, without claiming a fixed host port. @@ -211,7 +213,13 @@ def terminate(_signum, _frame): assert updated['public_url'] == origin, 'The new public URL did not take effect' assert any(p['id'] == project_data['id'] for p in get('/core/v1/projects')['data']), 'Project was lost' assert get('/v1/files/' + uploaded['id'], headers=api)['bytes'] == len(content), 'Uploaded file metadata was lost' - assert 'Bundled node installation metadata verified' not in private_logs(key, project_key), 'Completed initialization recopied metadata' + assert '"stage":"metadata_verification"' not in private_logs(key, project_key), 'Completed initialization recopied metadata' + compose('stop', 'core') + stopped_logs = private_logs(key, project_key) + assert any( + '"msg":"Stage completed"' in line and '"stage":"shutdown"' in line for line in stopped_logs.splitlines() + ), 'Normal stop did not complete the shutdown stage' + assert not any('"msg":"Stage failed"' in line and '"stage":"running"' in line for line in stopped_logs.splitlines()), 'Normal stop was reported as a runtime failure' print('PASS: startup, origin validation, sign-in, API, upload, node installer and persistent installation', flush=True) except BaseException: # Service status identifies failed containers without dumping secret-bearing logs. diff --git a/services/core/cmd/oac/init.go b/services/core/cmd/oac/init.go index 63a8ca93..818fc59a 100644 --- a/services/core/cmd/oac/init.go +++ b/services/core/cmd/oac/init.go @@ -7,13 +7,13 @@ import ( "encoding/hex" "encoding/json" "errors" - "fmt" "io/fs" "os" "path/filepath" "strings" "syscall" + "github.com/MiniMax-AI/OpenAgentCore/internal/obs/log" "github.com/google/uuid" ) @@ -89,7 +89,7 @@ func readRelease(root string, release releaseIdentity) (map[string][]byte, error if manifest.SourceCommit != release.revision || manifest.Platform != "linux/amd64" { return nil, errors.New("release identity mismatch") } - fmt.Println("Bundled node installation metadata verified") + log.Bg().Debug("Bundled node installation metadata verified") return files, nil } @@ -107,7 +107,11 @@ func writeOwned(path string, data []byte) error { if err := chown(temporary, 65532, 65532); err != nil { return err } - return os.Rename(temporary, path) + if err := os.Rename(temporary, path); err != nil { + return err + } + log.Bg().Debug("Installation file prepared", "path", path, "uid", 65532, "mode", "0600") + return nil } func ownedDir(path string, mode os.FileMode, uid int) error { @@ -134,7 +138,10 @@ type installReceipt struct { Files map[string]string `json:"files"` } -func initialize(root string, release releaseIdentity, fetch func() (map[string][]byte, error)) error { +func initialize(root string, release releaseIdentity, fetch func() (map[string][]byte, error)) (err error) { + finish := log.StartStage("directories", "component", "init", "data_dir", root, "revision", release.revision) + next := func(stage string) { finish(nil); finish = log.StartStage(stage, "component", "init") } + defer func() { finish(err) }() if err := os.Chmod(root, 0o755); err != nil { return err } @@ -148,6 +155,7 @@ func initialize(root string, release releaseIdentity, fetch func() (map[string][ return err } } + next("installation_lock") lock, err := os.OpenFile(filepath.Join(root, "secrets", ".init.lock"), os.O_CREATE|os.O_WRONLY, 0o600) if err != nil { return err @@ -156,6 +164,7 @@ func initialize(root string, release releaseIdentity, fetch func() (map[string][ if err := syscall.Flock(int(lock.Fd()), syscall.LOCK_EX); err != nil { return err } + next("installation_check") marker := filepath.Join(root, "installation.json") if raw, err := os.ReadFile(marker); err == nil { var receipt installReceipt @@ -174,7 +183,7 @@ func initialize(root string, release releaseIdentity, fetch func() (map[string][ return errors.New("installation files changed; restore the matching data directory") } } - fmt.Println("Existing installation verified") + log.Bg().Info("Existing installation verified") return nil } else if !errors.Is(err, fs.ErrNotExist) { return err @@ -188,10 +197,12 @@ func initialize(root string, release releaseIdentity, fetch func() (map[string][ return errors.New("existing data requires its original installation files") } } + next("metadata_verification") files, err := fetch() if err != nil { return err } + next("metadata_publication") prefix := "node-payload/releases/" + release.revision + "/" names := []string{} for _, name := range releaseMembers { @@ -215,6 +226,7 @@ func initialize(root string, release releaseIdentity, fetch func() (map[string][ }); err != nil { return err } + next("secrets") generators := []struct { name string generate func() string @@ -250,6 +262,7 @@ func initialize(root string, release releaseIdentity, fetch func() (map[string][ return err } names = append(names, "secrets/core/core-key-digests.json", "node-payload/active.json") + next("installation_receipt") receipt := installReceipt{SourceCommit: release.revision, Files: map[string]string{}} for _, name := range names { if receipt.Files[name], err = fileDigest(filepath.Join(root, name)); err != nil { @@ -260,7 +273,7 @@ func initialize(root string, release releaseIdentity, fetch func() (map[string][ if err := writeOwned(marker, raw); err != nil { return err } - fmt.Println("Installation initialized; print the sign-in key with: docker compose exec web oac-web core-key") + log.Bg().Info("Installation initialized", "next_action", "docker compose exec web oac-web core-key") return nil } diff --git a/services/core/cmd/oac/main.go b/services/core/cmd/oac/main.go index 254993a5..b3ffc3c3 100644 --- a/services/core/cmd/oac/main.go +++ b/services/core/cmd/oac/main.go @@ -13,6 +13,8 @@ import ( "path/filepath" "strings" "syscall" + + "github.com/MiniMax-AI/OpenAgentCore/internal/obs/log" ) func main() { @@ -38,6 +40,7 @@ func usage() { func run(ctx context.Context, command string, args []string) error { switch command { case "init": + log.Init(log.ConfigFromEnv()) return initCommand() } root, err := installDir() diff --git a/services/core/cmd/server/configuration_diagnostics.go b/services/core/cmd/server/configuration_diagnostics.go new file mode 100644 index 00000000..18a97aab --- /dev/null +++ b/services/core/cmd/server/configuration_diagnostics.go @@ -0,0 +1,22 @@ +package main + +import ( + "errors" + "fmt" + + "github.com/MiniMax-AI/OpenAgentCore/internal/obs/log" +) + +// reason must be authored validation text, never a parser or dependency error. +func configurationFailure(setting, reason string, cause error) error { + fields := []any{"setting", setting, "reason", reason} + if cause != nil { + fields = append(fields, log.ErrorFields(cause)...) + } + log.Bg().Error("Core configuration invalid", fields...) + message := setting + ": " + reason + if cause != nil { + return fmt.Errorf("%s: %w", message, cause) + } + return errors.New(message) +} diff --git a/services/core/cmd/server/configuration_diagnostics_test.go b/services/core/cmd/server/configuration_diagnostics_test.go new file mode 100644 index 00000000..25369bfa --- /dev/null +++ b/services/core/cmd/server/configuration_diagnostics_test.go @@ -0,0 +1,75 @@ +package main + +import ( + "bytes" + "errors" + "log/slog" + "os" + "path/filepath" + "strings" + "testing" + + "github.com/MiniMax-AI/OpenAgentCore/internal/obs/log" +) + +func TestConfigurationDiagnosticsPreserveFactsWithoutSecrets(t *testing.T) { + tests := []struct { + setting, content, reason string + read func() error + }{ + {"OAC_CREDENTIAL_KEY_FILE", "private-secret", "base64-encoded", func() error { _, err := credentialCipher(); return err }}, + {"OAC_CORE_KEY_DIGESTS_FILE", `["private-secret"]`, "unique lowercase SHA-256", func() error { _, err := deploymentAdminAuthenticator(); return err }}, + {"OAC_HISTORY_SETTINGS_FILE", `{"endpoint":"https://user:private-secret@example.test/metrics","transport":"otlp_http"}`, "canonical absolute", func() error { _, err := loadRuntimeHistoryConfig(); return err }}, + } + for _, test := range tests { + t.Run(test.setting, func(t *testing.T) { + var output bytes.Buffer + previous := slog.Default() + slog.SetDefault(slog.New(slog.NewJSONHandler(&output, nil))) + t.Cleanup(func() { slog.SetDefault(previous) }) + path := filepath.Join(t.TempDir(), "settings") + t.Setenv(test.setting, path) + if err := test.read(); !errors.Is(err, os.ErrNotExist) { + t.Fatalf("file cause was discarded: %v", err) + } + logged := output.String() + if !strings.Contains(logged, test.setting) || !strings.Contains(logged, path) || !strings.Contains(logged, `"error_kind":"not_found"`) { + t.Fatalf("missing actionable file diagnostic: %s", logged) + } + output.Reset() + if err := os.WriteFile(path, []byte(test.content), 0600); err != nil { + t.Fatal(err) + } + if err := test.read(); err == nil { + t.Fatal("invalid file accepted") + } + logged = output.String() + if !strings.Contains(logged, test.setting) || !strings.Contains(logged, test.reason) || strings.Contains(logged, "private-secret") { + t.Fatalf("unsafe or incomplete validation diagnostic: %s", logged) + } + }) + } +} + +func TestCoreRejectsInvalidDatabaseConfigurationBeforeMigrations(t *testing.T) { + log.Init(log.ConfigFromEnv()) + for _, setting := range []string{"OAC_PUBLIC_URL", "OAC_INSTALLATION_ID_FILE", "OAC_EXECUTION_CONCURRENCY", "OAC_DEFAULT_HARNESS", "OAC_HARNESSES", "OAC_WRITE_AUDIT_RETENTION", "OAC_OAUTH_TRUSTED_ORIGINS", "OAC_LOG_LEVEL", "OAC_LOG_FORMAT", "OAC_LOG_ADD_SOURCE", "OAC_DATABASE_PASSWORD_FILE"} { + t.Setenv(setting, "") + } + for _, databaseURL := range []string{"", "postgres://user:private-secret@%zz/db"} { + t.Run(databaseURL[:min(len(databaseURL), 8)], func(t *testing.T) { + var output bytes.Buffer + previous := slog.Default() + slog.SetDefault(slog.New(slog.NewJSONHandler(&output, nil))) + t.Cleanup(func() { slog.SetDefault(previous) }) + t.Setenv("OAC_DATABASE_URL", databaseURL) + if err := run(); err == nil { + t.Fatal("invalid database configuration accepted") + } + logged := output.String() + if !strings.Contains(logged, `"setting":"OAC_DATABASE_URL"`) || !strings.Contains(logged, `"msg":"Stage failed"`) || !strings.Contains(logged, `"stage":"configuration"`) || strings.Contains(logged, "database_migrations") || strings.Contains(logged, "private-secret") { + t.Fatalf("unsafe or misleading configuration diagnostic: %s", logged) + } + }) + } +} diff --git a/services/core/cmd/server/credential_cipher.go b/services/core/cmd/server/credential_cipher.go index 60a6754c..fa532b47 100644 --- a/services/core/cmd/server/credential_cipher.go +++ b/services/core/cmd/server/credential_cipher.go @@ -2,7 +2,6 @@ package main import ( "encoding/base64" - "errors" "os" "strings" @@ -16,11 +15,11 @@ func credentialCipher() (*credentialcrypto.Cipher, error) { } content, err := os.ReadFile(path) if err != nil { - return nil, errors.New("cannot read OAC_CREDENTIAL_KEY_FILE") + return nil, configurationFailure("OAC_CREDENTIAL_KEY_FILE", "cannot read credential key file", err) } key, err := base64.StdEncoding.Strict().DecodeString(strings.TrimSpace(string(content))) if err != nil || len(key) != 32 { - return nil, errors.New("OAC_CREDENTIAL_KEY_FILE must contain a base64-encoded random 32-byte key") + return nil, configurationFailure("OAC_CREDENTIAL_KEY_FILE", "must contain a base64-encoded random 32-byte key", nil) } return credentialcrypto.New(key) } diff --git a/services/core/cmd/server/main.go b/services/core/cmd/server/main.go index 04430569..403e0182 100644 --- a/services/core/cmd/server/main.go +++ b/services/core/cmd/server/main.go @@ -25,6 +25,7 @@ import ( "context" "errors" "fmt" + "net" "net/http" "os" "os/signal" @@ -81,16 +82,19 @@ func main() { return } if err := run(); err != nil { - log.Bg().Error("oac-core startup failed", "error", err) os.Exit(1) } } -func run() error { +func run() (runErr error) { + log.Init(log.ConfigFromEnv()) + finish := log.StartStage("configuration", "component", "core", "revision", buildRevision) + next := func(stage string) { finish(nil); finish = log.StartStage(stage, "component", "core") } + defer func() { finish(runErr) }() if err := processconfig.Check(); err != nil { + log.Bg().Error("Core configuration invalid", "error", err) return err } - log.Init(log.ConfigFromEnv()) public, err := processconfig.PublicURL() if err != nil { return err @@ -102,10 +106,14 @@ func run() error { logConfigurationSources() databaseURL, err := databaseurl.FromEnvironment() if err != nil { - return err + return configurationFailure("OAC_DATABASE_URL / OAC_DATABASE_PASSWORD_FILE", err.Error(), nil) } if databaseURL == "" { - return errors.New("OAC_DATABASE_URL is required") + return configurationFailure("OAC_DATABASE_URL", "is required", nil) + } + databaseConfig, err := pgxpool.ParseConfig(databaseURL) + if err != nil { + return configurationFailure("OAC_DATABASE_URL", "invalid Agents API database configuration", nil) } credentialKey, err := credentialCipher() if err != nil { @@ -113,22 +121,25 @@ func run() error { } ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM) defer stop() + next("database_migrations") migrating, cancelMigration := context.WithTimeout(ctx, 2*time.Minute) err = migrations.Apply(migrating, databaseURL) cancelMigration() if err != nil { return fmt.Errorf("Agents API database migration failed: %w", err) } - pool, err := pgxpool.New(ctx, databaseURL) + next("database_connection") + pool, err := pgxpool.NewWithConfig(ctx, databaseConfig) if err != nil { - return errors.New("invalid Agents API database configuration") + return configurationFailure("OAC_DATABASE_URL", "invalid Agents API database configuration", nil) } defer pool.Close() ready, cancel := context.WithTimeout(ctx, 10*time.Second) defer cancel() if err := pool.Ping(ready); err != nil { - return errors.New("Agents API database connection failed") + return fmt.Errorf("Agents API database connection failed: %w", err) } + next("services") engine, err := processconfig.DefaultHarness() if err != nil { return err @@ -213,6 +224,7 @@ func run() error { defer func() { cancelAuditCleanup(); <-auditCleanupDone }() var workerDone chan error var worker *execution.Worker + next("runtime_setup") managedNodes, err := configureManagedNodes(deploymentService, deploymentStore, sandboxProviders, public, func(ctx context.Context) error { if worker == nil { return errors.New("sandbox execution owner is unavailable") @@ -312,8 +324,12 @@ func run() error { Sessions: sessionService, SessionsReader: sessionStore, ManagedRuntimes: managed, MaxConcurrentExecutions: concurrency} + next("execution_worker") lease, err := pgunit.AcquireLease(ctx, pool) if err != nil { + if errors.Is(err, pgunit.ErrLeaseHeld) { + log.Bg().Error("Core execution lease unavailable", "reason", "another Core execution service owns this database") + } return err } deploymentExecution, err := deployment.NewExecutionOperations(deploymentService, deploymentpg.NewExecution(lease, credentialKey)) @@ -359,7 +375,7 @@ func run() error { _, sampleErr = deploymentStore.SampleHostHistory(sampleCtx) } cancel() - if !result.Complete { + if !result.Complete && sampleErr == nil { sampleErr = errors.New("incomplete Runtime sampling sweep") } metrics.ReportJob("runtime_sampler", result.CompletedAt, metricPtr(int64(result.Observed)), metricPtr(int64(result.Failed)), sampleErr) @@ -367,7 +383,7 @@ func run() error { if result.Complete { log.Bg().Debug("Runtime history sampling sweep complete", fields...) } else { - log.Bg().Warn("Runtime history sampling sweep incomplete", fields...) + log.Bg().Warn("Runtime history sampling sweep incomplete", append(fields, log.ErrorFields(sampleErr)...)...) } }, }) @@ -434,6 +450,7 @@ func run() error { ConfigurationDiscovery: managedNodes.setup, } } + next("http_listener") handler, err := api.NewHandler(deps) if err != nil { return err @@ -449,19 +466,33 @@ func run() error { } addr := serverAddress() server := &http.Server{Addr: addr, Handler: handler, ReadHeaderTimeout: 10 * time.Second, ReadTimeout: 30 * time.Second, WriteTimeout: 30 * time.Second, IdleTimeout: 60 * time.Second} + listener, err := net.Listen("tcp", addr) + if err != nil { + return err + } + log.Bg().Info("Core HTTP listener ready", "address", listener.Addr().String()) + next("running") done := make(chan error, 1) - go func() { done <- server.ListenAndServe() }() + go func() { done <- server.Serve(listener) }() select { case err := <-done: return err case err := <-workerDone: workerDone = nil + normalStop := ctx.Err() != nil && errors.Is(err, context.Canceled) + if normalStop { + next("shutdown") + } stop() shutdown, cancel := context.WithTimeout(context.Background(), 10*time.Second) defer cancel() - _ = server.Shutdown(shutdown) + shutdownErr := server.Shutdown(shutdown) + if normalStop { + return shutdownErr + } return err case <-ctx.Done(): + next("shutdown") shutdown, cancel := context.WithTimeout(context.Background(), 10*time.Second) defer cancel() return server.Shutdown(shutdown) diff --git a/services/core/cmd/server/managed_nodes.go b/services/core/cmd/server/managed_nodes.go index 8e9add39..3dd0e09f 100644 --- a/services/core/cmd/server/managed_nodes.go +++ b/services/core/cmd/server/managed_nodes.go @@ -34,7 +34,7 @@ func configureManagedNodes(nodes *deployment.Service, reader deployment.Reader, return nil, err } if publicURL == "" { - return nil, errors.New("OAC_INSTALLATION_ID_FILE requires OAC_PUBLIC_URL, the origin nodes and sandboxes use to reach Core") + return nil, configurationFailure("OAC_PUBLIC_URL", "is required for nodes and sandboxes to reach Core", nil) } closeProvider := func() {} result := &managedNodes{closeProvider: closeProvider} @@ -49,7 +49,7 @@ func configureManagedNodes(nodes *deployment.Service, reader deployment.Reader, return nil, err } if result.admin == nil { - return nil, errors.New("Web sandbox setup requires OAC_CORE_KEY_DIGESTS_FILE with the Core key digest") + return nil, configurationFailure("OAC_CORE_KEY_DIGESTS_FILE", "is required for Web sandbox setup", nil) } result.hub = node.NewHub(node.HubOptions{ Generations: func(ctx context.Context, n node.Identity, connection string, epoch uint64, health node.Health) error { @@ -117,13 +117,17 @@ func deploymentAdminAuthenticator() (*api.DeploymentAuthenticator, error) { } raw, err := os.ReadFile(path) if err != nil { - return nil, errors.New("cannot read OAC_CORE_KEY_DIGESTS_FILE") + return nil, configurationFailure("OAC_CORE_KEY_DIGESTS_FILE", "cannot read Core key digests file", err) } var digests []string if json.Unmarshal(raw, &digests) != nil || len(digests) == 0 { - return nil, errors.New("OAC_CORE_KEY_DIGESTS_FILE must contain a JSON array of Core key SHA-256 digests") + return nil, configurationFailure("OAC_CORE_KEY_DIGESTS_FILE", "must contain a JSON array of Core key SHA-256 digests", nil) } - return api.NewDeploymentAuthenticator(digests) + authenticator, err := api.NewDeploymentAuthenticator(digests) + if err != nil { + return nil, configurationFailure("OAC_CORE_KEY_DIGESTS_FILE", "must contain unique lowercase SHA-256 digests", nil) + } + return authenticator, nil } func serverAddress() string { diff --git a/services/core/cmd/server/runtime_history.go b/services/core/cmd/server/runtime_history.go index 308e84a5..813782e4 100644 --- a/services/core/cmd/server/runtime_history.go +++ b/services/core/cmd/server/runtime_history.go @@ -11,6 +11,7 @@ import ( "strings" "time" + "github.com/MiniMax-AI/OpenAgentCore/internal/obs/log" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/coremetrics" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimehistory" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimehistory/postgresreader" @@ -93,16 +94,16 @@ func loadRuntimeHistoryConfig() (runtimeHistoryConfig, error) { if file := os.Getenv("OAC_HISTORY_SETTINGS_FILE"); file != "" { raw, err := os.ReadFile(file) if err != nil { - return config, errors.New("cannot read OAC_HISTORY_SETTINGS_FILE") + return config, configurationFailure("OAC_HISTORY_SETTINGS_FILE", "cannot read Runtime history settings file", err) } decoder := json.NewDecoder(bytes.NewReader(raw)) decoder.DisallowUnknownFields() if decoder.Decode(&config) != nil || decoder.Decode(new(any)) != io.EOF { - return config, errors.New("invalid Runtime history configuration") + return config, configurationFailure("OAC_HISTORY_SETTINGS_FILE", "invalid Runtime history configuration", nil) } } if err := validateRuntimeHistoryConfig(config); err != nil { - return config, err + return config, configurationFailure("OAC_HISTORY_SETTINGS_FILE", err.Error(), nil) } if config.QueueCapacity == 0 { config.QueueCapacity = defaultRuntimeHistoryQueueCapacity @@ -133,6 +134,9 @@ func runHistoryCleanup(ctx context.Context, prune func(context.Context) (int64, count, err := prune(pruneCtx) cancel() reportCleanupResult(metrics, "history_cleanup", count, err) + if err != nil && ctx.Err() == nil { + log.Ctx(ctx).Warn("Runtime history retention cleanup failed", log.ErrorFields(err)...) + } select { case <-ctx.Done(): return diff --git a/services/core/cmd/server/write_audit.go b/services/core/cmd/server/write_audit.go index 48153352..ce6d7162 100644 --- a/services/core/cmd/server/write_audit.go +++ b/services/core/cmd/server/write_audit.go @@ -38,7 +38,7 @@ func runWriteAuditCleanup(ctx context.Context, s writeAuditPruner, retention tim cancel() reportCleanupResult(metrics, "audit_cleanup", count, err) if err != nil && ctx.Err() == nil { - log.Ctx(ctx).Warn("Write audit retention cleanup failed") + log.Ctx(ctx).Warn("Write audit retention cleanup failed", log.ErrorFields(err)...) } select { case <-ctx.Done(): diff --git a/services/core/internal/processconfig/config.go b/services/core/internal/processconfig/config.go index c7f1805b..7f4ccd17 100644 --- a/services/core/internal/processconfig/config.go +++ b/services/core/internal/processconfig/config.go @@ -1,10 +1,12 @@ // Package processconfig is the process settings Core loads from its environment. // Startup and `oac-core check-config` both call Check. Settings reports the -// effective values for GET /core/v1/installation. Errors name the variable and -// never include its value. +// effective values for GET /core/v1/installation. Validation errors name the +// variable without echoing settings or file contents; file errors retain paths +// and their typed filesystem causes. package processconfig import ( + "fmt" "os" "slices" "strconv" @@ -121,7 +123,7 @@ func InstallationID() (string, error) { } raw, err := os.ReadFile(path) if err != nil { - return "", configErr("OAC_INSTALLATION_ID_FILE must name a readable file") + return "", fmt.Errorf("OAC_INSTALLATION_ID_FILE must name a readable file: %w", err) } value := strings.TrimSpace(string(raw)) if id, err := uuid.Parse(value); err != nil || id == uuid.Nil || id.String() != value { diff --git a/services/core/internal/processconfig/config_test.go b/services/core/internal/processconfig/config_test.go index 9f8dfa53..1047c2d4 100644 --- a/services/core/internal/processconfig/config_test.go +++ b/services/core/internal/processconfig/config_test.go @@ -1,6 +1,8 @@ package processconfig import ( + "errors" + "os" "strings" "testing" ) @@ -67,3 +69,11 @@ func TestHarnessesDefaultToEveryQualifiedHarness(t *testing.T) { t.Fatal("unqualified harness enabled") } } + +func TestInstallationIDReadErrorRetainsItsCause(t *testing.T) { + t.Setenv("OAC_INSTALLATION_ID_FILE", "/missing-oac-installation/installation.id") + _, err := InstallationID() + if !errors.Is(err, os.ErrNotExist) || !strings.Contains(err.Error(), "OAC_INSTALLATION_ID_FILE") { + t.Fatalf("err = %v", err) + } +} diff --git a/services/core/tests/official_diagnostics.py b/services/core/tests/official_diagnostics.py index 31a794f7..62c6282d 100644 --- a/services/core/tests/official_diagnostics.py +++ b/services/core/tests/official_diagnostics.py @@ -1,21 +1,18 @@ """Bounded failure facts for the real-service official-client fixture.""" import json -import re import subprocess import sys EVENTS = { "Core process configuration loaded from the process environment": "service_start", - "oac-core startup failed": "service_exit", -} -ERRORS = { - "context canceled": "context_canceled", - "context deadline exceeded": "deadline_exceeded", - "conn closed": "connection_closed", - "invalid input": "invalid_input", + "Stage failed": "service_exit", } +ERRORS = {name: name for name in ( + "canceled", "timeout", "not_found", "permission_denied", "connection_refused", + "connection_reset", "unexpected_eof", "eof", "database_error", +)} SQLSTATES = { "08006": "connection_failure", "22P02": "invalid_text_representation", "23503": "foreign_key_violation", "23505": "unique_violation", @@ -42,12 +39,15 @@ def failure_facts(output, sensitive): omitted += 1 continue event = {"event": EVENTS[entry["msg"]]} - error = entry.get("error") + if entry["msg"] == "Stage failed" and entry.get("component") != "core": + omitted += 1 + continue + error = entry.get("error_kind") if isinstance(error, str): event["error"] = ERRORS.get(error, "unclassified") - match = re.search(r"\(SQLSTATE ([0-9A-Z]{5})\)$", error) - if match and match[1] in SQLSTATES: - event.update(sqlstate=match[1], error=SQLSTATES[match[1]]) + state = entry.get("sqlstate") + if isinstance(state, str) and state in SQLSTATES: + event.update(sqlstate=state, error=SQLSTATES[state]) events.append(event) return {"events": events, "omitted_lines": omitted} diff --git a/services/core/tests/official_diagnostics_test.py b/services/core/tests/official_diagnostics_test.py index a8fe1b35..454a8557 100644 --- a/services/core/tests/official_diagnostics_test.py +++ b/services/core/tests/official_diagnostics_test.py @@ -14,7 +14,7 @@ class DiagnosticsTests(unittest.TestCase): def test_failed_owned_process_reports_exit_and_safe_error(self): with tempfile.TemporaryFile(mode="w+") as log: - code = 'import json; print(json.dumps({"msg":"oac-core startup failed","error":"ERROR: deadlock detected (SQLSTATE 40P01)"})); raise SystemExit(7)' + code = 'import json; print(json.dumps({"msg":"Stage failed","component":"core","error_kind":"database_error","sqlstate":"40P01"})); raise SystemExit(7)' process = subprocess.Popen([sys.executable, "-c", code], stdout=log, stderr=log) self.assertEqual(process.wait(timeout=10), 7) diagnostics = io.StringIO() @@ -31,8 +31,8 @@ def test_failed_owned_process_reports_exit_and_safe_error(self): def test_secrets_and_untrusted_fields_never_enter_diagnostics(self): secret = "dynamic-issued-project-key" - lines = [json.dumps({"msg": "oac-core startup failed", "error": secret}), - json.dumps({"msg": "oac-core startup failed", "error": "postgres://user:unknown-password@host/db", "Authorization": "Bearer unknown-key", "body": "private request"}), + lines = [json.dumps({"msg": "Stage failed", "component": "core", "error_kind": secret}), + json.dumps({"msg": "Stage failed", "component": "core", "error_kind": "postgres://user:unknown-password@host/db", "Authorization": "Bearer unknown-key", "body": "private request"}), "panic: private request", json.dumps({"msg": ["invalid log"]})] diagnostics = io.StringIO() primary = RuntimeError("primary") @@ -68,7 +68,7 @@ def test_success_keeps_secret_assertion_and_has_no_diagnostic(self): output = io.StringIO() finish_server(None, io.StringIO(""), [], None, output) self.assertEqual(output.getvalue(), "") - self.assertEqual(len(failure_facts('\n'.join([json.dumps({"msg": "oac-core startup failed"})] * 100), [])['events']), 32) + self.assertEqual(len(failure_facts('\n'.join([json.dumps({"msg": "Stage failed", "component": "core"})] * 100), [])['events']), 32) if __name__ == "__main__":