From b789925ae5c55171a585f6b9fc68937f663e34a1 Mon Sep 17 00:00:00 2001 From: npc0-hue Date: Fri, 7 Aug 2026 12:01:26 +0800 Subject: [PATCH] fix: recover run runtime state and logs --- .../tasks.md | 28 +++---- platform/api/resource_handlers.go | 19 +++++ platform/domain/control.go | 3 + platform/domain/log_ingest.go | 18 +++++ platform/domain/resources.go | 29 +++---- platform/dto/control.go | 6 ++ platform/dto/log_ingest.go | 24 ++++++ platform/dto/resources.go | 2 + platform/model/resources.go | 75 +++++++++++-------- platform/protocol/run-contracts.md | 3 +- platform/service/control_test.go | 25 +++++++ platform/service/log_ingest.go | 30 ++++++++ platform/service/log_ingest_test.go | 20 +++++ platform/service/resources.go | 1 + .../service/server_lifecycle_projection.go | 22 ++++++ platform/validator/job_channel.go | 6 ++ platform/validator/log_ingest.go | 9 +++ platform_web/api/types.ts | 1 + .../components/ServerLiveOperations.tsx | 6 +- .../contracts/serverManagement.test.ts | 12 ++- platform_web/contracts/serverManagement.ts | 8 ++ platform_web/contracts/workspace.ts | 7 +- platform_web/pages/ServersPage.tsx | 21 +++--- 23 files changed, 297 insertions(+), 78 deletions(-) diff --git a/openspec/changes/repair-run-runtime-state-and-log-recovery/tasks.md b/openspec/changes/repair-run-runtime-state-and-log-recovery/tasks.md index acaa96e..c27ba08 100644 --- a/openspec/changes/repair-run-runtime-state-and-log-recovery/tasks.md +++ b/openspec/changes/repair-run-runtime-state-and-log-recovery/tasks.md @@ -8,29 +8,29 @@ ## 1. Platform Contracts And Projections -- [ ] 1.1 Add typed domain, DTO, API, validation, and protocol contracts for signed Run log-stream progress queries scoped to the authenticated Run endpoint and bound server instance. -- [ ] 1.2 Implement Platform log-stream progress lookup that returns only the latest acknowledged sequence and cannot disclose log bodies, host paths, credentials, or another server's stream metadata. -- [ ] 1.3 Extend Run lifecycle observations with stable managed-process identity and ordering data, then make lifecycle projection idempotent and reject stale state regressions. -- [ ] 1.4 Project autonomous Run recovered and exit observations through the existing signed lifecycle channel, preserving requested-stop versus unexpected-exit classification. +- [x] 1.1 Add typed domain, DTO, API, validation, and protocol contracts for signed Run log-stream progress queries scoped to the authenticated Run endpoint and bound server instance. +- [x] 1.2 Implement Platform log-stream progress lookup that returns only the latest acknowledged sequence and cannot disclose log bodies, host paths, credentials, or another server's stream metadata. +- [x] 1.3 Extend Run lifecycle observations with stable managed-process identity and ordering data, then make lifecycle projection idempotent and reject stale state regressions. +- [x] 1.4 Project autonomous Run recovered and exit observations through the existing signed lifecycle channel, preserving requested-stop versus unexpected-exit classification. - [ ] 1.5 Expose a server runtime observation view that combines the persisted lifecycle projection with generic bound-endpoint heartbeat freshness without changing lifecycle state solely because Run is unavailable. -- [ ] 1.6 Add focused Platform tests for report ordering/idempotency, authorization and scope of log progress, progress values after durable ingest, and fresh versus unverified runtime observation. +- [x] 1.6 Add focused Platform tests for report ordering/idempotency, authorization and scope of log progress, progress values after durable ingest, and fresh versus unverified runtime observation. ## 2. Independent Run Recovery -- [ ] 2.1 Update the shared/copyable Run-Platform protocol types and API client in `git@git.npc0.com:admin343/run.git` for lifecycle observation ordering and signed log-stream progress reconciliation. -- [ ] 2.2 Add atomic per-stream allocated and acknowledged watermark persistence to the Run log spool, including restart loading and acknowledgement-before-segment-deletion ordering. -- [ ] 2.3 Replace the worker-global in-memory log counter with stream-specific allocation restored from the spool watermark and pending durable segments. -- [ ] 2.4 Reconcile signed Platform stream progress before a stable Run-bound stream with no local watermark emits new entries; cover newly created and recreated-spool cases. -- [ ] 2.5 Classify acknowledged-range conflicts and sequence gaps as durable recovery failures, quarantine the affected spool segment with redacted diagnostics, and resume only after safe watermark reconciliation. -- [ ] 2.6 Make autonomous process supervision report observed exit and startup-recovery transitions through the lifecycle channel, with retry-safe process identity and ordering metadata. -- [ ] 2.7 Define and test graceful Run shutdown behavior that preserves durable state and never reports a server stop unless its generic supervisor observed that process state. +- [x] 2.1 Update the shared/copyable Run-Platform protocol types and API client in `git@git.npc0.com:admin343/run.git` for lifecycle observation ordering and signed log-stream progress reconciliation. +- [x] 2.2 Add atomic per-stream allocated and acknowledged watermark persistence to the Run log spool, including restart loading and acknowledgement-before-segment-deletion ordering. +- [x] 2.3 Replace the worker-global in-memory log counter with stream-specific allocation restored from the spool watermark and pending durable segments. +- [x] 2.4 Reconcile signed Platform stream progress before a stable Run-bound stream with no local watermark emits new entries; cover newly created and recreated-spool cases. +- [x] 2.5 Classify acknowledged-range conflicts and sequence gaps as durable recovery failures, quarantine the affected spool segment with redacted diagnostics, and resume only after safe watermark reconciliation. +- [x] 2.6 Make autonomous process supervision report observed exit and startup-recovery transitions through the lifecycle channel, with retry-safe process identity and ordering metadata. +- [x] 2.7 Define and test graceful Run shutdown behavior that preserves durable state and never reports a server stop unless its generic supervisor observed that process state. - [ ] 2.8 Add Run unit tests for per-stream interleaving, restart continuity, missing-watermark progress lookup, conflict quarantine, process exit reporting, and Windows supervisor recovery. ## 3. Management Runtime Presentation -- [ ] 3.1 Extend Platform Web API types and server-management contracts to consume lifecycle projection and runtime observation freshness separately. +- [x] 3.1 Extend Platform Web API types and server-management contracts to consume lifecycle projection and runtime observation freshness separately. - [ ] 3.2 Update server list and server detail status UI so stale `running` is presented as last observed with a Run offline/unverified qualifier, not confirmed online. -- [ ] 3.3 Update the management terminal header and empty/error states to show that live output awaits Run recovery while preserving accepted bounded SSE history. +- [x] 3.3 Update the management terminal header and empty/error states to show that live output awaits Run recovery while preserving accepted bounded SSE history. - [ ] 3.4 Add focused frontend tests for fresh, stale, offline, and recovered Run observations plus terminal presentation during log recovery. ## 4. Cross-Repository Verification And Release diff --git a/platform/api/resource_handlers.go b/platform/api/resource_handlers.go index 1e78ed8..30f2666 100644 --- a/platform/api/resource_handlers.go +++ b/platform/api/resource_handlers.go @@ -137,6 +137,7 @@ func (h *coreHandlers) register(mux *http.ServeMux) { mux.HandleFunc("/api/v1/run/jobs/cancel", h.requireRunSignature(h.runJobCancelPoll)) mux.HandleFunc("/api/v1/run/jobs/reconcile", h.requireRunSignature(h.runJobReconcile)) mux.HandleFunc("/api/v1/run/logs/batches", h.requireRunSignature(h.runLogBatchIngest)) + mux.HandleFunc("/api/v1/run/logs/progress", h.requireRunSignature(h.runLogStreamProgress)) mux.HandleFunc("/api/v1/run/artifacts/open", h.requireRunSignature(h.runArtifactOpen)) mux.HandleFunc("/api/v1/run/artifacts/chunks", h.requireRunSignature(h.runArtifactChunkUpload)) mux.HandleFunc("/api/v1/run/artifacts/status", h.requireRunSignature(h.runArtifactStatus)) @@ -2037,6 +2038,24 @@ func (h *coreHandlers) runLogBatchIngest(w http.ResponseWriter, r *http.Request) writeJSON(w, http.StatusOK, dto.LogBatchIngestFromDomain(result)) } +func (h *coreHandlers) runLogStreamProgress(w http.ResponseWriter, r *http.Request) { + if r.Method != http.MethodPost { + writeMethodNotAllowed(w, http.MethodPost) + return + } + request, err := decodeJSON[dto.RunLogStreamProgressRequest](r) + if err != nil { + writeDecodeError(w, err) + return + } + result, err := h.core.GetRunLogStreamProgress(request.ToDomain()) + if err != nil { + writeServiceError(w, err) + return + } + writeJSON(w, http.StatusOK, dto.RunLogStreamProgressFromDomain(result)) +} + // runArtifactOpen godoc // @Summary Open run artifact upload transfer // @Description Lets a registered run endpoint open a resumable upload transfer for a scoped artifact owner. diff --git a/platform/domain/control.go b/platform/domain/control.go index 186cefa..50cec3b 100644 --- a/platform/domain/control.go +++ b/platform/domain/control.go @@ -62,6 +62,9 @@ type RunLifecycleReport struct { Progress RunJobProgressReport Message string ErrorCode string + ManagedProcessID string + ObservationSeq uint64 + ObservedAt time.Time ExecutionResult JobExecutionResult } diff --git a/platform/domain/log_ingest.go b/platform/domain/log_ingest.go index 660b351..3ae9078 100644 --- a/platform/domain/log_ingest.go +++ b/platform/domain/log_ingest.go @@ -48,6 +48,24 @@ type LogStreamCursorResult struct { LatestSeq uint64 } +// RunLogStreamProgress is a signed Run-only request for sequence recovery. +// It deliberately exposes no log body or host-local information. +type RunLogStreamProgress struct { + RunEndpointID string + SessionToken string + ServerInstanceID string + LogStreamID string +} + +type RunLogStreamProgressResult struct { + Accepted bool + RunEndpointID string + ServerInstanceID string + LogStreamID string + LatestSeq uint64 + ServerTime time.Time +} + type LogStreamEvent struct { ServerInstanceID string Stream LogStream diff --git a/platform/domain/resources.go b/platform/domain/resources.go index d5cc23a..dad3757 100644 --- a/platform/domain/resources.go +++ b/platform/domain/resources.go @@ -777,19 +777,22 @@ type ServerInstance struct { // DeploymentTargetID identifies an optional operator-selected deployment // target for post-creation deployment operations. It never selects a // distribution builder or replaces the server's generated Run endpoint. - DeploymentTargetID string - RunEndpointID string - Name string - OwnerUserID string - AdminUserIDs []string - State ServerInstanceState - ConfigVersion int - ConfigKey string - ConfigContent string - ConfigChecksum string - ConfigUpdatedAt time.Time - CreatedAt time.Time - UpdatedAt time.Time + DeploymentTargetID string + RunEndpointID string + Name string + OwnerUserID string + AdminUserIDs []string + State ServerInstanceState + LifecycleProcessID string + LifecycleObservationSeq uint64 + LifecycleObservedAt time.Time + ConfigVersion int + ConfigKey string + ConfigContent string + ConfigChecksum string + ConfigUpdatedAt time.Time + CreatedAt time.Time + UpdatedAt time.Time // Deployment stores operator-supplied deployment inputs. Its protected path // and command values are never included in normal platform read projections. Deployment ServerDeploymentDefinition diff --git a/platform/dto/control.go b/platform/dto/control.go index 68ee8c1..716b676 100644 --- a/platform/dto/control.go +++ b/platform/dto/control.go @@ -66,6 +66,9 @@ type RunLifecycleReportRequest struct { Progress JobProgressBody `json:"progress"` Message string `json:"message,omitempty"` ErrorCode string `json:"errorCode,omitempty"` + ManagedProcessID string `json:"managedProcessId,omitempty"` + ObservationSeq uint64 `json:"observationSeq,omitempty"` + ObservedAt time.Time `json:"observedAt,omitempty"` ExecutionResult RunJobExecutionResultBody `json:"executionResult,omitempty"` } @@ -122,6 +125,9 @@ func (request RunLifecycleReportRequest) ToDomain() domain.RunLifecycleReport { Progress: progressReportToDomain(request.Progress), Message: request.Message, ErrorCode: request.ErrorCode, + ManagedProcessID: request.ManagedProcessID, + ObservationSeq: request.ObservationSeq, + ObservedAt: request.ObservedAt, ExecutionResult: domain.JobExecutionResult{Kind: request.ExecutionResult.Kind, ProcessState: request.ExecutionResult.ProcessState, ExitClassification: request.ExecutionResult.ExitClassification, ExitCode: request.ExecutionResult.ExitCode, Version: request.ExecutionResult.Version, Checksum: request.ExecutionResult.Checksum, SizeBytes: request.ExecutionResult.SizeBytes, AuditSummary: request.ExecutionResult.AuditSummary, Content: request.ExecutionResult.Content, ServerDeploymentEvidence: serverDeploymentEvidenceToDomain(request.ExecutionResult.ServerDeploymentEvidence), DeploymentReceipt: deploymentReceiptToDomain(request.ExecutionResult.DeploymentReceipt)}, } } diff --git a/platform/dto/log_ingest.go b/platform/dto/log_ingest.go index 87fa908..f3e736e 100644 --- a/platform/dto/log_ingest.go +++ b/platform/dto/log_ingest.go @@ -53,6 +53,22 @@ type LogStreamCursorResponse struct { LatestSeq uint64 `json:"latestSeq"` } +type RunLogStreamProgressRequest struct { + RunEndpointID string `json:"runEndpointId"` + SessionToken string `json:"sessionToken"` + ServerInstanceID string `json:"serverInstanceId"` + LogStreamID string `json:"logStreamId"` +} + +type RunLogStreamProgressResponse struct { + Accepted bool `json:"accepted"` + RunEndpointID string `json:"runEndpointId"` + ServerInstanceID string `json:"serverInstanceId"` + LogStreamID string `json:"logStreamId"` + LatestSeq uint64 `json:"latestSeq"` + ServerTime time.Time `json:"serverTime"` +} + type LogStreamEventResponse struct { ServerInstanceID string `json:"serverInstanceId"` StreamID string `json:"streamId"` @@ -92,6 +108,10 @@ func (request LogStreamCursorRequest) ToDomain() domain.LogStreamCursorQuery { } } +func (request RunLogStreamProgressRequest) ToDomain() domain.RunLogStreamProgress { + return domain.RunLogStreamProgress{RunEndpointID: request.RunEndpointID, SessionToken: request.SessionToken, ServerInstanceID: request.ServerInstanceID, LogStreamID: request.LogStreamID} +} + func LogBatchIngestFromDomain(result domain.LogBatchIngestResult) LogBatchIngestResponse { return LogBatchIngestResponse{ Accepted: result.Accepted, @@ -114,6 +134,10 @@ func LogStreamCursorFromDomain(result domain.LogStreamCursorResult) LogStreamCur } } +func RunLogStreamProgressFromDomain(result domain.RunLogStreamProgressResult) RunLogStreamProgressResponse { + return RunLogStreamProgressResponse{Accepted: result.Accepted, RunEndpointID: result.RunEndpointID, ServerInstanceID: result.ServerInstanceID, LogStreamID: result.LogStreamID, LatestSeq: result.LatestSeq, ServerTime: result.ServerTime} +} + func LogStreamEventFromDomain(event domain.LogStreamEvent) LogStreamEventResponse { event = domain.CopyLogStreamEvent(event) return LogStreamEventResponse{ diff --git a/platform/dto/resources.go b/platform/dto/resources.go index e8aa7a8..7b03df8 100644 --- a/platform/dto/resources.go +++ b/platform/dto/resources.go @@ -587,6 +587,7 @@ type ServerInstanceResponse struct { OwnerUserID string `json:"ownerUserId,omitempty"` AdminUserIDs []string `json:"adminUserIds"` State domain.ServerInstanceState `json:"state"` + LifecycleObservedAt *time.Time `json:"lifecycleObservedAt,omitempty"` ConfigVersion int `json:"configVersion"` ConfigKey string `json:"configKey,omitempty"` ConfigChecksum string `json:"configChecksum,omitempty"` @@ -1640,6 +1641,7 @@ func ServerInstanceFromDomain(instance domain.ServerInstance) ServerInstanceResp OwnerUserID: instance.OwnerUserID, AdminUserIDs: adminUserIDs, State: instance.State, + LifecycleObservedAt: optionalTime(instance.LifecycleObservedAt), ConfigVersion: instance.ConfigVersion, ConfigKey: instance.ConfigKey, ConfigChecksum: instance.ConfigChecksum, diff --git a/platform/model/resources.go b/platform/model/resources.go index fc00ec8..419a44f 100644 --- a/platform/model/resources.go +++ b/platform/model/resources.go @@ -216,7 +216,10 @@ type ServerInstance struct { // AdminUserIDs identifies server-scoped administrator accounts. AdminUserIDs []string `json:"adminUserIds" db:"admin_user_ids"` // State is the server lifecycle state. - State domain.ServerInstanceState `json:"state" db:"state"` + State domain.ServerInstanceState `json:"state" db:"state"` + LifecycleProcessID string `json:"lifecycleProcessId,omitempty" db:"lifecycle_process_id"` + LifecycleObservationSeq uint64 `json:"lifecycleObservationSeq,omitempty" db:"lifecycle_observation_seq"` + LifecycleObservedAt time.Time `json:"lifecycleObservedAt,omitempty" db:"lifecycle_observed_at"` // ConfigVersion is the platform-managed optimistic concurrency version. ConfigVersion int `json:"configVersion" db:"config_version"` // ConfigKey is the logical configuration target, never a host path. @@ -739,43 +742,49 @@ func remoteAccessFromDomain(remote domain.GamePluginRemoteAccess) GamePluginRemo func ServerInstanceFromDomain(instance domain.ServerInstance) ServerInstance { return ServerInstance{ - ID: instance.ID, - PluginID: instance.PluginID, - PluginVersion: instance.PluginVersion, - DeploymentTargetID: instance.DeploymentTargetID, - RunEndpointID: instance.RunEndpointID, - Name: instance.Name, - OwnerUserID: instance.OwnerUserID, - AdminUserIDs: domain.CopyStringSlice(instance.AdminUserIDs), - State: instance.State, - ConfigVersion: instance.ConfigVersion, - ConfigKey: instance.ConfigKey, - ConfigContent: instance.ConfigContent, - ConfigChecksum: instance.ConfigChecksum, - ConfigUpdatedAt: instance.ConfigUpdatedAt, - CreatedAt: instance.CreatedAt, - UpdatedAt: instance.UpdatedAt, + ID: instance.ID, + PluginID: instance.PluginID, + PluginVersion: instance.PluginVersion, + DeploymentTargetID: instance.DeploymentTargetID, + RunEndpointID: instance.RunEndpointID, + Name: instance.Name, + OwnerUserID: instance.OwnerUserID, + AdminUserIDs: domain.CopyStringSlice(instance.AdminUserIDs), + State: instance.State, + LifecycleProcessID: instance.LifecycleProcessID, + LifecycleObservationSeq: instance.LifecycleObservationSeq, + LifecycleObservedAt: instance.LifecycleObservedAt, + ConfigVersion: instance.ConfigVersion, + ConfigKey: instance.ConfigKey, + ConfigContent: instance.ConfigContent, + ConfigChecksum: instance.ConfigChecksum, + ConfigUpdatedAt: instance.ConfigUpdatedAt, + CreatedAt: instance.CreatedAt, + UpdatedAt: instance.UpdatedAt, } } func (instance ServerInstance) ToDomain() domain.ServerInstance { return domain.ServerInstance{ - ID: instance.ID, - PluginID: instance.PluginID, - PluginVersion: instance.PluginVersion, - DeploymentTargetID: instance.DeploymentTargetID, - RunEndpointID: instance.RunEndpointID, - Name: instance.Name, - OwnerUserID: instance.OwnerUserID, - AdminUserIDs: domain.CopyStringSlice(instance.AdminUserIDs), - State: instance.State, - ConfigVersion: instance.ConfigVersion, - ConfigKey: instance.ConfigKey, - ConfigContent: instance.ConfigContent, - ConfigChecksum: instance.ConfigChecksum, - ConfigUpdatedAt: instance.ConfigUpdatedAt, - CreatedAt: instance.CreatedAt, - UpdatedAt: instance.UpdatedAt, + ID: instance.ID, + PluginID: instance.PluginID, + PluginVersion: instance.PluginVersion, + DeploymentTargetID: instance.DeploymentTargetID, + RunEndpointID: instance.RunEndpointID, + Name: instance.Name, + OwnerUserID: instance.OwnerUserID, + AdminUserIDs: domain.CopyStringSlice(instance.AdminUserIDs), + State: instance.State, + LifecycleProcessID: instance.LifecycleProcessID, + LifecycleObservationSeq: instance.LifecycleObservationSeq, + LifecycleObservedAt: instance.LifecycleObservedAt, + ConfigVersion: instance.ConfigVersion, + ConfigKey: instance.ConfigKey, + ConfigContent: instance.ConfigContent, + ConfigChecksum: instance.ConfigChecksum, + ConfigUpdatedAt: instance.ConfigUpdatedAt, + CreatedAt: instance.CreatedAt, + UpdatedAt: instance.UpdatedAt, } } diff --git a/platform/protocol/run-contracts.md b/platform/protocol/run-contracts.md index f0ce1a7..2370c1f 100644 --- a/platform/protocol/run-contracts.md +++ b/platform/protocol/run-contracts.md @@ -63,7 +63,7 @@ Platform-owned Run distribution builds embed an autonomous lifecycle plan for th The plan is build input for the generated package, not a machine-side job-channel payload. Generated Run registration must not be treated as a trigger to enqueue `process.start`, `process.install`, or `process.status` work; Platform state converges from Run heartbeats, logs, lifecycle reports, supervised process facts, and terminal job/report messages. Platform and Run must not add game-specific hardcoding to interpret the plan. -Autonomous lifecycle reports use `POST /api/v1/run/lifecycle/report` with the active Run session and signed envelope when required. The route accepts only bounded terminal lifecycle facts for `process.install`, `process.start`, `process.stop`, or `process.status`; it validates the server/run binding, records audit evidence, and projects server state from Run-reported process facts without creating or completing a Platform job. +Autonomous lifecycle reports use `POST /api/v1/run/lifecycle/report` with the active Run session and signed envelope when required. The route accepts only bounded terminal lifecycle facts for `process.install`, `process.start`, `process.stop`, or `process.status`; it validates the server/run binding, records audit evidence, and projects server state from Run-reported process facts without creating or completing a Platform job. A managed-process report includes an opaque `managedProcessId`, monotonic `observationSeq`, and `observedAt`; retries are idempotent and a lower sequence cannot regress a newer fact for that process. ## Log Ingest @@ -77,6 +77,7 @@ Named log DTOs: - `LogBatchIngestRequest` - `LogBatchIngestResponse` +- `RunLogStreamProgressRequest` / `RunLogStreamProgressResponse`: signed Run-only sequence recovery for a server-bound `run...*` stream. The response contains only the latest acknowledged sequence. - `LogEntry` - `LogStreamCursorRequest` - `LogStreamCursorResponse` diff --git a/platform/service/control_test.go b/platform/service/control_test.go index bbffc63..ad11af4 100644 --- a/platform/service/control_test.go +++ b/platform/service/control_test.go @@ -4,6 +4,7 @@ import ( "errors" "strings" "testing" + "time" "browser.local/platform/domain" "browser.local/platform/repo" @@ -464,6 +465,30 @@ func TestCoreServiceRunLifecycleReportProjectsGeneratedRunFacts(t *testing.T) { } } +func TestCoreServiceRejectsStaleManagedProcessObservation(t *testing.T) { + svc := newTestCoreService() + plugin := createGeneratedRunStatusPlugin(t, svc) + instance := domain.ServerInstance{ID: "managed-observation-order", PluginID: plugin.ID, PluginVersion: plugin.Version, RunEndpointID: dedicatedRunEndpointID("managed-observation-order"), Name: "Managed Observation Order", State: domain.ServerInstanceStateDraft, ConfigVersion: 1} + if err := svc.store.ServerInstances().Create(instance); err != nil { + t.Fatalf("create server: %v", err) + } + registered := registerGeneratedRunForStatusTest(t, svc, instance, plugin.ID) + observedAt := time.Date(2026, 8, 7, 10, 0, 0, 0, time.UTC) + report := func(sequence uint64, state string, classification string) { + t.Helper() + if _, err := svc.ReportRunLifecycle(domain.RunLifecycleReport{RunEndpointID: instance.RunEndpointID, SessionToken: registered.SessionToken, ServerInstanceID: instance.ID, Capability: domain.LifecycleCapabilityStatus, State: domain.JobStateSucceeded, Progress: domain.RunJobProgressReport{Percent: 100}, ManagedProcessID: "sha256:managed-process", ObservationSeq: sequence, ObservedAt: observedAt.Add(time.Duration(sequence) * time.Second), ExecutionResult: domain.JobExecutionResult{Kind: "process", ProcessState: state, ExitClassification: classification}}); err != nil { + t.Fatalf("report sequence %d: %v", sequence, err) + } + } + report(1, "running", "") + report(2, "exited", "unexpected-exit") + report(1, "running", "") + stored, err := svc.GetServerInstance(instance.ID) + if err != nil || stored.State != domain.ServerInstanceStateFailed || stored.LifecycleObservationSeq != 2 { + t.Fatalf("stale running observation must not regress exit projection: server=%+v err=%v", stored, err) + } +} + func TestCoreServiceGeneratedRunRegistrationDoesNotDispatchStatusReconciliation(t *testing.T) { svc := newTestCoreService() plugin := createGeneratedRunStatusPlugin(t, svc) diff --git a/platform/service/log_ingest.go b/platform/service/log_ingest.go index 95b7095..78fd3a9 100644 --- a/platform/service/log_ingest.go +++ b/platform/service/log_ingest.go @@ -219,6 +219,36 @@ func (svc *CoreService) QueryLogStream(query domain.LogStreamCursorQuery) (domai }), nil } +func (svc *CoreService) GetRunLogStreamProgress(request domain.RunLogStreamProgress) (domain.RunLogStreamProgressResult, error) { + if err := validator.ValidateRunLogStreamProgress(request); err != nil { + return domain.RunLogStreamProgressResult{}, err + } + if _, err := svc.validatedRunSession(request.RunEndpointID, request.SessionToken); err != nil { + return domain.RunLogStreamProgressResult{}, err + } + instance, err := svc.store.ServerInstances().Get(request.ServerInstanceID) + if err != nil { + return domain.RunLogStreamProgressResult{}, err + } + if instance.RunEndpointID != request.RunEndpointID { + return domain.RunLogStreamProgressResult{}, validationError("runEndpointId must match server instance") + } + if !strings.HasPrefix(request.LogStreamID, "run."+request.RunEndpointID+"."+request.ServerInstanceID+".") { + return domain.RunLogStreamProgressResult{}, validationError("logStreamId is not a bound run stream") + } + latest := uint64(0) + stream, err := svc.store.LogStreams().Get(request.LogStreamID) + if err == nil { + if stream.ServerInstanceID != instance.ID { + return domain.RunLogStreamProgressResult{}, validationError("logStreamId does not belong to server instance") + } + latest = stream.LatestSeq + } else if !errors.Is(err, repo.ErrNotFound) { + return domain.RunLogStreamProgressResult{}, err + } + return domain.RunLogStreamProgressResult{Accepted: true, RunEndpointID: request.RunEndpointID, ServerInstanceID: instance.ID, LogStreamID: request.LogStreamID, LatestSeq: latest, ServerTime: svc.now()}, nil +} + func validateLogBatchStream(batch domain.LogBatchIngest, stream domain.LogStream) error { if stream.ID != batch.LogStreamID { return validationError("logStreamId must match stream") diff --git a/platform/service/log_ingest_test.go b/platform/service/log_ingest_test.go index d2d3d3b..3033969 100644 --- a/platform/service/log_ingest_test.go +++ b/platform/service/log_ingest_test.go @@ -39,6 +39,26 @@ func TestCoreServiceIngestsLogBatchAndQueriesCursor(t *testing.T) { } } +func TestCoreServiceReturnsBoundRunLogStreamProgress(t *testing.T) { + svc, sessionToken := newRegisteredLogIngestService(t) + streamID := "run.run-local.server-1.stdout" + if _, err := svc.CreateLogStream(domain.LogStream{ID: streamID, ServerInstanceID: "server-1", Source: domain.LogStreamSourceProcess, StreamKey: "stdout", StorageBackend: domain.LogStorageBackendLocalSegments, RetentionPolicy: "default"}); err != nil { + t.Fatalf("create run stream: %v", err) + } + batch := validLogBatch(t, sessionToken, 1, 2) + batch.LogStreamID = streamID + if _, err := svc.IngestLogBatch(batch); err != nil { + t.Fatalf("ingest run stream: %v", err) + } + progress, err := svc.GetRunLogStreamProgress(domain.RunLogStreamProgress{RunEndpointID: "run-local", SessionToken: sessionToken, ServerInstanceID: "server-1", LogStreamID: streamID}) + if err != nil || !progress.Accepted || progress.LatestSeq != 2 { + t.Fatalf("unexpected progress: %+v err=%v", progress, err) + } + if _, err := svc.GetRunLogStreamProgress(domain.RunLogStreamProgress{RunEndpointID: "run-local", SessionToken: sessionToken, ServerInstanceID: "server-1", LogStreamID: "run.run-local.other.stdout"}); err == nil { + t.Fatal("expected progress scope validation failure") + } +} + func TestCoreServicePublishesLogEventsForAcceptedBatch(t *testing.T) { svc, sessionToken := newRegisteredLogIngestService(t) createLogStreamFixture(t, svc) diff --git a/platform/service/resources.go b/platform/service/resources.go index 505aa41..0ea7dc1 100644 --- a/platform/service/resources.go +++ b/platform/service/resources.go @@ -209,6 +209,7 @@ type Core interface { SubscribeLogEvents(string) (LogEventSubscription, error) SubscribeLogEventsForSession(string, string) (LogEventSubscription, error) IngestLogBatch(domain.LogBatchIngest) (domain.LogBatchIngestResult, error) + GetRunLogStreamProgress(domain.RunLogStreamProgress) (domain.RunLogStreamProgressResult, error) QueryLogStream(domain.LogStreamCursorQuery) (domain.LogStreamCursorResult, error) ListGamePlayersForSession(string, domain.GamePlayerFilter) ([]domain.GamePlayer, error) GetGamePlayerProfileForSession(string, string) (domain.GamePlayerProfile, error) diff --git a/platform/service/server_lifecycle_projection.go b/platform/service/server_lifecycle_projection.go index 8b4e0da..c41783d 100644 --- a/platform/service/server_lifecycle_projection.go +++ b/platform/service/server_lifecycle_projection.go @@ -29,8 +29,20 @@ func (svc *CoreService) ReportRunLifecycle(report domain.RunLifecycleReport) (do stamp := svc.now() nextState, projected := lifecycleProjectedState(report.Capability, report.State, report.ExecutionResult) + if lifecycleObservationIsStale(instance, report) { + projected = false + nextState = instance.State + } if projected { instance.State = nextState + if report.ManagedProcessID != "" { + instance.LifecycleProcessID = report.ManagedProcessID + instance.LifecycleObservationSeq = report.ObservationSeq + instance.LifecycleObservedAt = report.ObservedAt + if instance.LifecycleObservedAt.IsZero() { + instance.LifecycleObservedAt = stamp + } + } instance.UpdatedAt = stamp if err := validator.ValidateServerInstance(instance); err != nil { return domain.RunLifecycleReportResult{}, err @@ -49,6 +61,16 @@ func (svc *CoreService) ReportRunLifecycle(report domain.RunLifecycleReport) (do return domain.CopyRunLifecycleReportResult(domain.RunLifecycleReportResult{Accepted: true, RunEndpointID: report.RunEndpointID, ServerInstanceID: report.ServerInstanceID, ProjectedState: nextState, ServerTime: stamp}), nil } +func lifecycleObservationIsStale(instance domain.ServerInstance, report domain.RunLifecycleReport) bool { + if report.ManagedProcessID == "" || instance.LifecycleProcessID == "" { + return false + } + if report.ManagedProcessID == instance.LifecycleProcessID { + return report.ObservationSeq <= instance.LifecycleObservationSeq + } + return !report.ObservedAt.IsZero() && !instance.LifecycleObservedAt.IsZero() && report.ObservedAt.Before(instance.LifecycleObservedAt) +} + func lifecycleReportSummary(report domain.RunLifecycleReport, projectedState domain.ServerInstanceState, projected bool) string { for _, candidate := range []string{report.ExecutionResult.AuditSummary, report.Progress.Message, report.Message, report.ErrorCode} { if strings.TrimSpace(candidate) != "" { diff --git a/platform/validator/job_channel.go b/platform/validator/job_channel.go index c57082d..c2ea5ed 100644 --- a/platform/validator/job_channel.go +++ b/platform/validator/job_channel.go @@ -69,6 +69,12 @@ func ValidateRunLifecycleReport(report domain.RunLifecycleReport) error { violations = appendProgressViolations(violations, report.Progress) violations = appendMessageLength(violations, "message", report.Message) violations = appendMessageLength(violations, "errorCode", report.ErrorCode) + if report.ManagedProcessID != "" && report.ObservationSeq == 0 { + violations = append(violations, "observationSeq is required with managedProcessId") + } + if report.ObservationSeq > 0 && report.ManagedProcessID == "" { + violations = append(violations, "managedProcessId is required with observationSeq") + } if len([]byte(report.ExecutionResult.Content)) > maxJobChannelMessageLength*256 { violations = append(violations, "executionResult.content is too large") } diff --git a/platform/validator/log_ingest.go b/platform/validator/log_ingest.go index fd05995..042b43f 100644 --- a/platform/validator/log_ingest.go +++ b/platform/validator/log_ingest.go @@ -87,6 +87,15 @@ func ValidateLogStreamCursorQuery(query domain.LogStreamCursorQuery) error { return finish(violations) } +func ValidateRunLogStreamProgress(request domain.RunLogStreamProgress) error { + var violations []string + violations = appendRequired(violations, "runEndpointId", request.RunEndpointID) + violations = appendRequired(violations, "sessionToken", request.SessionToken) + violations = appendRequired(violations, "serverInstanceId", request.ServerInstanceID) + violations = appendRequired(violations, "logStreamId", request.LogStreamID) + return finish(violations) +} + func LogEntriesChecksum(entries []domain.LogEntry) (string, error) { stable := make([]logEntryChecksumBody, len(entries)) for i, entry := range entries { diff --git a/platform_web/api/types.ts b/platform_web/api/types.ts index abc0590..4179574 100644 --- a/platform_web/api/types.ts +++ b/platform_web/api/types.ts @@ -462,6 +462,7 @@ export interface ServerInstanceResponse { ownerUserId?: string; adminUserIds: string[]; state: ServerInstanceState; + lifecycleObservedAt?: string; configVersion: number; configKey?: string; configChecksum?: string; diff --git a/platform_web/components/ServerLiveOperations.tsx b/platform_web/components/ServerLiveOperations.tsx index 1f2afee..c23692c 100644 --- a/platform_web/components/ServerLiveOperations.tsx +++ b/platform_web/components/ServerLiveOperations.tsx @@ -314,7 +314,7 @@ export function ServerManagementTerminalDrawer({ open, serverId, serverName, plu function handleTerminalScroll() { const output = outputRef.current; - if (!output) return; + if (!output || initialHistoryPendingRef.current) return; const nextFollowLatest = output.scrollHeight - output.clientHeight - output.scrollTop <= 24; followLatestRef.current = nextFollowLatest; setFollowLatest(nextFollowLatest); @@ -369,7 +369,7 @@ export function ServerManagementTerminalDrawer({ open, serverId, serverName, plu
{serverName} - 最近历史 + SSE 实时推送 · {streams.status === "ready" ? `${terminalStreams.length} 个日志源` : streams.status === "loading" ? "读取日志源" : "日志源异常"} · {followLatest ? "自动置底" : "已解锁滚动"} + 已接受历史 + SSE 实时推送 · 实时输出等待 Run 连通与水位恢复 · {streams.status === "ready" ? `${terminalStreams.length} 个日志源` : streams.status === "loading" ? "读取日志源" : "日志源异常"} · {followLatest ? "自动置底" : "已解锁滚动"}
@@ -378,7 +378,7 @@ export function ServerManagementTerminalDrawer({ open, serverId, serverName, plu
{streams.status === "error" &&
LOGS{streams.reason}
} - {streams.status === "ready" && streams.data.length === 0 &&
LOGS暂无日志源。需要 Run 上报或历史日志回填后,这里才会持续追加。
} + {streams.status === "ready" && streams.data.length === 0 &&
LOGS暂无已接受日志。Run 恢复连接并完成日志水位校准后,新输出会继续追加。
} {lines.map((line) =>
{line.streamKey || line.level || "LOG"}{line.text}
)}
diff --git a/platform_web/contracts/serverManagement.test.ts b/platform_web/contracts/serverManagement.test.ts index 58e426b..037aea5 100644 --- a/platform_web/contracts/serverManagement.test.ts +++ b/platform_web/contracts/serverManagement.test.ts @@ -1,7 +1,7 @@ import { describe, expect, it } from "vitest"; -import { canStartServer, canStopServer } from "./serverManagement"; -import type { ServerInstanceState } from "../api/types"; +import { canStartServer, canStopServer, runtimeObservationFreshness } from "./serverManagement"; +import type { RunEndpointResponse, ServerInstanceResponse, ServerInstanceState } from "../api/types"; describe("server management lifecycle contracts", () => { it("allows explicit starts from recoverable non-running states", () => { @@ -16,4 +16,12 @@ describe("server management lifecycle contracts", () => { expect(canStopServer("running")).toBe(true); expect(canStopServer("failed")).toBe(false); }); + + it("distinguishes a fresh Run observation from an unverified historical lifecycle state", () => { + const instance = { id: "server-1", state: "running", runEndpointId: "run-1" } as ServerInstanceResponse; + const endpoint = { id: "run-1", status: "online", lastHeartbeatAt: "2026-08-07T10:00:00Z" } as RunEndpointResponse; + expect(runtimeObservationFreshness(instance, endpoint, Date.parse("2026-08-07T10:00:30Z"))).toBe("fresh"); + expect(runtimeObservationFreshness(instance, endpoint, Date.parse("2026-08-07T10:01:00Z"))).toBe("unverified"); + expect(runtimeObservationFreshness(instance, { ...endpoint, status: "offline" }, Date.parse("2026-08-07T10:00:01Z"))).toBe("unverified"); + }); }); diff --git a/platform_web/contracts/serverManagement.ts b/platform_web/contracts/serverManagement.ts index bb41b9a..b21f91b 100644 --- a/platform_web/contracts/serverManagement.ts +++ b/platform_web/contracts/serverManagement.ts @@ -61,6 +61,14 @@ export interface ServerManagementSummary { failed: number; } +export type RuntimeObservationFreshness = "fresh" | "unverified"; + +export function runtimeObservationFreshness(instance: ServerInstanceResponse, endpoint: RunEndpointResponse | undefined, now = Date.now()): RuntimeObservationFreshness { + if (!endpoint || endpoint.status !== "online") return "unverified"; + const heartbeat = Date.parse(endpoint.lastHeartbeatAt); + return Number.isFinite(heartbeat) && now-heartbeat <= 45_000 ? "fresh" : "unverified"; +} + export const emptyServerCreateForm: ServerCreateFormState = { id: "", name: "", diff --git a/platform_web/contracts/workspace.ts b/platform_web/contracts/workspace.ts index 7dad641..fca7c7f 100644 --- a/platform_web/contracts/workspace.ts +++ b/platform_web/contracts/workspace.ts @@ -1,4 +1,4 @@ -import type { JobResponse, ServerInstanceResponse, ServerMetricsResponse, UserContactProfile, UserStatus, UserThemePreferenceResponse } from "../api/types"; +import type { JobResponse, RunEndpointResponse, ServerInstanceResponse, ServerMetricsResponse, UserContactProfile, UserStatus, UserThemePreferenceResponse } from "../api/types"; export type WorkspaceRole = "platformAdmin" | "serverOwner" | "serverAdmin"; @@ -89,8 +89,9 @@ export interface OperationRecord { } export interface ServerCardView { - instance: ServerInstanceResponse; - metrics?: ServerMetricsResponse; + instance: ServerInstanceResponse; + endpoint?: RunEndpointResponse; + metrics?: ServerMetricsResponse; pendingJobs: number; activeJobs?: number; failedJobs?: number; diff --git a/platform_web/pages/ServersPage.tsx b/platform_web/pages/ServersPage.tsx index 6dad0bc..c9d8452 100644 --- a/platform_web/pages/ServersPage.tsx +++ b/platform_web/pages/ServersPage.tsx @@ -23,8 +23,9 @@ import type { PageComponentProps } from "../contracts/page"; import { canDeleteServer, defaultServerCreateForm, - endpointLabel, - pluginCreateInputDefaults, + endpointLabel, + pluginCreateInputDefaults, + runtimeObservationFreshness, type ServerCreateFormState } from "../contracts/serverManagement"; import { summarizeServerOperations } from "../contracts/operationsConsole"; @@ -152,15 +153,16 @@ export function ServersPage({ session, operations, onNavigate }: PageComponentPr const cards = useMemo( () => - summarizeServerOperations(instances, metrics, jobs).map((summary) => ({ - instance: summary.instance, - metrics: summary.metrics, + summarizeServerOperations(instances, metrics, jobs).map((summary) => ({ + instance: summary.instance, + endpoint: endpoints.find((endpoint) => endpoint.id === summary.instance.runEndpointId), + metrics: summary.metrics, pendingJobs: summary.activeJobs, activeJobs: summary.activeJobs, failedJobs: summary.failedJobs, latestJob: summary.latestJob })), - [instances, jobs, metrics] + [endpoints, instances, jobs, metrics] ); const visibleCards = useMemo(() => filterServerCards(cards, keyword, statusFilter), [cards, keyword, statusFilter]); @@ -725,8 +727,9 @@ interface ServerCardProps { } function ServerCard({ card, metricsPending, metricsUnavailable, canManage, deleteDisabledReason, onOpen, onEdit, onQuickAction, onDelete }: ServerCardProps) { - const { instance, metrics, pendingJobs, failedJobs = 0 } = card; - const online = serverIsOnline(instance.state); + const { instance, endpoint, metrics, pendingJobs, failedJobs = 0 } = card; + const online = serverIsOnline(instance.state); + const freshness = runtimeObservationFreshness(instance, endpoint); const canDelete = deleteDisabledReason === ""; const canOpenActions = canManage || canDelete; const metricsWaiting = metrics?.source === "run-metrics-pending"; @@ -815,7 +818,7 @@ function ServerCard({ card, metricsPending, metricsUnavailable, canManage, delet {instance.name} {instance.id} - {stateLabel(instance.state)} + {freshness === "fresh" ? stateLabel(instance.state) : `最后观测:${stateLabel(instance.state)}(Run 未验证)`}