Fix run log ingest and terminal console layout
This commit is contained in:
@@ -171,9 +171,10 @@ const (
|
||||
type LogStreamSource string
|
||||
|
||||
const (
|
||||
LogStreamSourceProcess LogStreamSource = "process"
|
||||
LogStreamSourceFile LogStreamSource = "file"
|
||||
LogStreamSourcePlugin LogStreamSource = "plugin"
|
||||
LogStreamSourceProcess LogStreamSource = "process"
|
||||
LogStreamSourceFile LogStreamSource = "file"
|
||||
LogStreamSourcePlugin LogStreamSource = "plugin"
|
||||
LogStreamSourceManagementProgram LogStreamSource = "management-program"
|
||||
)
|
||||
|
||||
type LogStorageBackend string
|
||||
|
||||
@@ -31,7 +31,7 @@ func (svc *CoreService) IngestLogBatch(batch domain.LogBatchIngest) (domain.LogB
|
||||
if err != nil {
|
||||
return domain.LogBatchIngestResult{}, err
|
||||
}
|
||||
if exists && record.LastSeq == batch.LastSeq && record.Checksum == batch.Checksum {
|
||||
if exists && record.LastSeq == batch.LastSeq && logBatchRecordMatches(record, batch) {
|
||||
if err := svc.projectGamePlayerEvents(projectionBatch); err != nil {
|
||||
return domain.LogBatchIngestResult{}, err
|
||||
}
|
||||
@@ -50,7 +50,7 @@ func (svc *CoreService) IngestLogBatch(batch domain.LogBatchIngest) (domain.LogB
|
||||
}
|
||||
return domain.LogBatchIngestResult{}, validationError("log batch conflicts with acknowledged range")
|
||||
}
|
||||
if batch.FirstSeq != stream.LatestSeq+1 {
|
||||
if stream.LatestSeq > 0 && batch.FirstSeq != stream.LatestSeq+1 {
|
||||
return domain.LogBatchIngestResult{}, validationError("log batch firstSeq must follow latest acknowledged sequence")
|
||||
}
|
||||
|
||||
@@ -86,6 +86,16 @@ func (svc *CoreService) IngestLogBatch(batch domain.LogBatchIngest) (domain.LogB
|
||||
}, nil
|
||||
}
|
||||
|
||||
func logBatchRecordMatches(record domain.LogBatchRecord, batch domain.LogBatchIngest) bool {
|
||||
if record.Checksum == batch.Checksum {
|
||||
return true
|
||||
}
|
||||
if len(record.Entries) == 1 && len(batch.Entries) == 1 {
|
||||
return batch.Checksum == validator.LogLineChecksum(record.Entries[0].Line)
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
// sanitizeGamePlayerNetworkFields removes raw network material before the durable log body is written.
|
||||
func sanitizeGamePlayerNetworkFields(batch *domain.LogBatchIngest) {
|
||||
for index := range batch.Entries {
|
||||
|
||||
@@ -56,16 +56,80 @@ func TestCoreServiceLogBatchDuplicateAck(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestCoreServiceAcceptsAutoCreatedRunJobLogStreams(t *testing.T) {
|
||||
svc, sessionToken := newRegisteredLogIngestService(t)
|
||||
job, err := svc.CreateJob(domain.Job{
|
||||
ID: "job-run-logs",
|
||||
ServerInstanceID: "server-1",
|
||||
RunEndpointID: "run-local",
|
||||
Capability: domain.LifecycleCapabilityStart,
|
||||
IdempotencyKey: "job-run-logs",
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("create job: %v", err)
|
||||
}
|
||||
streamID := jobLogStreamID(job.ID, "stderr")
|
||||
stream, err := svc.GetLogStream(streamID)
|
||||
if err != nil {
|
||||
t.Fatalf("get auto-created job log stream: %v", err)
|
||||
}
|
||||
if stream.StreamKey != "stderr" || stream.Source != domain.LogStreamSourceProcess {
|
||||
t.Fatalf("unexpected stream metadata: %+v", stream)
|
||||
}
|
||||
|
||||
entry := domain.LogEntry{Seq: 2, Timestamp: time.Date(2026, 7, 3, 12, 0, 2, 0, time.UTC), Level: "info", Line: "stderr:server-ready"}
|
||||
batch := domain.LogBatchIngest{
|
||||
RunEndpointID: "run-local",
|
||||
SessionToken: sessionToken,
|
||||
LogStreamID: streamID,
|
||||
ServerInstanceID: "server-1",
|
||||
StreamKey: "stderr",
|
||||
Source: domain.LogStreamSourceProcess,
|
||||
FirstSeq: entry.Seq,
|
||||
LastSeq: entry.Seq,
|
||||
Compression: "none",
|
||||
Checksum: validator.LogLineChecksum(entry.Line),
|
||||
Entries: []domain.LogEntry{entry},
|
||||
}
|
||||
ack, err := svc.IngestLogBatch(batch)
|
||||
if err != nil {
|
||||
t.Fatalf("ingest run job log batch: %v", err)
|
||||
}
|
||||
if !ack.Accepted || ack.LatestSeq != entry.Seq {
|
||||
t.Fatalf("unexpected ack: %+v", ack)
|
||||
}
|
||||
duplicate, err := svc.IngestLogBatch(batch)
|
||||
if err != nil {
|
||||
t.Fatalf("ingest duplicate run job log batch: %v", err)
|
||||
}
|
||||
if !duplicate.Duplicate {
|
||||
t.Fatalf("expected duplicate ack, got %+v", duplicate)
|
||||
}
|
||||
query, err := svc.QueryLogStream(domain.LogStreamCursorQuery{LogStreamID: streamID, AfterSeq: 0, Limit: 10})
|
||||
if err != nil {
|
||||
t.Fatalf("query auto-created job log stream: %v", err)
|
||||
}
|
||||
if len(query.Entries) != 1 || query.Entries[0].Line != entry.Line || query.NextSeq != entry.Seq {
|
||||
t.Fatalf("unexpected query result: %+v", query)
|
||||
}
|
||||
}
|
||||
|
||||
func TestCoreServiceRejectsOutOfOrderAndConflictingLogBatches(t *testing.T) {
|
||||
svc, sessionToken := newRegisteredLogIngestService(t)
|
||||
createLogStreamFixture(t, svc)
|
||||
gap := validLogBatch(t, sessionToken, 2, 2)
|
||||
first := validLogBatch(t, sessionToken, 1, 1)
|
||||
if _, err := svc.IngestLogBatch(first); err != nil {
|
||||
t.Fatalf("ingest first batch: %v", err)
|
||||
}
|
||||
gap := validLogBatch(t, sessionToken, 3, 3)
|
||||
|
||||
_, err := svc.IngestLogBatch(gap)
|
||||
if err == nil || !strings.Contains(err.Error(), "firstSeq") {
|
||||
t.Fatalf("expected out-of-order rejection, got %v", err)
|
||||
}
|
||||
|
||||
svc, sessionToken = newRegisteredLogIngestService(t)
|
||||
createLogStreamFixture(t, svc)
|
||||
batch := validLogBatch(t, sessionToken, 1, 2)
|
||||
if _, err := svc.IngestLogBatch(batch); err != nil {
|
||||
t.Fatalf("ingest first batch: %v", err)
|
||||
@@ -139,7 +203,9 @@ func newRegisteredLogIngestService(t *testing.T) (*CoreService, string) {
|
||||
}); err != nil {
|
||||
t.Fatalf("create server instance: %v", err)
|
||||
}
|
||||
hello, err := svc.RegisterRunHello(validRunControlHello())
|
||||
helloRequest := validRunControlHello()
|
||||
helloRequest.CapabilityReport.Capabilities = append(helloRequest.CapabilityReport.Capabilities, domain.LifecycleCapabilityStart)
|
||||
hello, err := svc.RegisterRunHello(helloRequest)
|
||||
if err != nil {
|
||||
t.Fatalf("register run hello: %v", err)
|
||||
}
|
||||
|
||||
@@ -328,6 +328,9 @@ func NewCoreServiceWithDurableStores(store repo.Store, logStore LogBodyStore, ar
|
||||
for _, session := range sessions {
|
||||
service.artifactTransfers[session.TransferID] = domain.CopyArtifactTransferSession(session)
|
||||
}
|
||||
if err := service.recoverJobLogStreams(); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if err := service.recoverLogCursors(); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
@@ -343,6 +346,33 @@ func NewCoreServiceWithDurableStores(store repo.Store, logStore LogBodyStore, ar
|
||||
return service, nil
|
||||
}
|
||||
|
||||
func (svc *CoreService) recoverJobLogStreams() error {
|
||||
jobs, err := svc.store.Jobs().List(domain.JobFilter{})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
stamp := svc.now()
|
||||
for _, job := range jobs {
|
||||
if strings.TrimSpace(job.ID) == "" || strings.TrimSpace(job.ServerInstanceID) == "" {
|
||||
continue
|
||||
}
|
||||
instance, err := svc.store.ServerInstances().Get(job.ServerInstanceID)
|
||||
if errors.Is(err, repo.ErrNotFound) {
|
||||
continue
|
||||
}
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if instance.State == domain.ServerInstanceStateDeleted {
|
||||
continue
|
||||
}
|
||||
if err := svc.ensureJobLogStreams(job, stamp); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (svc *CoreService) recoverLogCursors() error {
|
||||
store, ok := svc.logStore.(interface{ LatestSeq(string) (uint64, error) })
|
||||
if !ok {
|
||||
@@ -2247,6 +2277,9 @@ func (svc *CoreService) CreateJob(job domain.Job) (domain.Job, error) {
|
||||
|
||||
existing, err := svc.store.Jobs().GetByIdempotency(job.RunEndpointID, job.IdempotencyKey)
|
||||
if err == nil {
|
||||
if err := svc.ensureJobLogStreams(existing, stamp); err != nil {
|
||||
return domain.Job{}, err
|
||||
}
|
||||
return existing, nil
|
||||
}
|
||||
if !errors.Is(err, repo.ErrNotFound) {
|
||||
@@ -2283,9 +2316,57 @@ func (svc *CoreService) CreateJob(job domain.Job) (domain.Job, error) {
|
||||
if err := svc.store.Jobs().Create(job); err != nil {
|
||||
return domain.Job{}, err
|
||||
}
|
||||
if err := svc.ensureJobLogStreams(job, stamp); err != nil {
|
||||
return domain.Job{}, err
|
||||
}
|
||||
return domain.CopyJob(job), nil
|
||||
}
|
||||
|
||||
func (svc *CoreService) ensureJobLogStreams(job domain.Job, stamp time.Time) error {
|
||||
if strings.TrimSpace(job.ServerInstanceID) == "" || strings.TrimSpace(job.ID) == "" {
|
||||
return nil
|
||||
}
|
||||
streams := []struct {
|
||||
key string
|
||||
source domain.LogStreamSource
|
||||
}{
|
||||
{key: "stdout", source: domain.LogStreamSourceProcess},
|
||||
{key: "stderr", source: domain.LogStreamSourceProcess},
|
||||
}
|
||||
if job.Capability == domain.JobCapabilityRemoteRunProgram {
|
||||
streams = append(streams,
|
||||
struct {
|
||||
key string
|
||||
source domain.LogStreamSource
|
||||
}{key: "management-program.stdout", source: domain.LogStreamSourceManagementProgram},
|
||||
struct {
|
||||
key string
|
||||
source domain.LogStreamSource
|
||||
}{key: "management-program.stderr", source: domain.LogStreamSourceManagementProgram},
|
||||
)
|
||||
}
|
||||
for _, item := range streams {
|
||||
stream := domain.LogStream{
|
||||
ID: jobLogStreamID(job.ID, item.key),
|
||||
ServerInstanceID: job.ServerInstanceID,
|
||||
Source: item.source,
|
||||
StreamKey: item.key,
|
||||
StorageBackend: domain.LogStorageBackendLocalSegments,
|
||||
RetentionPolicy: "default",
|
||||
CreatedAt: stamp,
|
||||
UpdatedAt: stamp,
|
||||
}
|
||||
if _, err := svc.CreateLogStream(stream); err != nil && !errors.Is(err, repo.ErrDuplicate) {
|
||||
return err
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func jobLogStreamID(jobID string, streamKey string) string {
|
||||
return fmt.Sprintf("job.%s.%s", jobID, streamKey)
|
||||
}
|
||||
|
||||
func (svc *CoreService) GetJob(id string) (domain.Job, error) {
|
||||
job, err := svc.store.Jobs().Get(id)
|
||||
if err != nil {
|
||||
|
||||
@@ -114,6 +114,10 @@ func TestCoreServiceCreateListGetWorkflows(t *testing.T) {
|
||||
if _, err := svc.GetJob(job.ID); err != nil {
|
||||
t.Fatalf("get job: %v", err)
|
||||
}
|
||||
streams, err := svc.ListLogStreams(domain.LogStreamFilter{ServerInstanceID: instance.ID})
|
||||
if err != nil || len(streams) != 2 {
|
||||
t.Fatalf("expected default job stdout/stderr streams, len=%d err=%v streams=%+v", len(streams), err, streams)
|
||||
}
|
||||
jobs, err := svc.ListJobs(domain.JobFilter{RunEndpointID: endpoint.ID})
|
||||
if err != nil || len(jobs) != 1 {
|
||||
t.Fatalf("list jobs: len=%d err=%v", len(jobs), err)
|
||||
@@ -154,8 +158,8 @@ func TestCoreServiceCreateListGetWorkflows(t *testing.T) {
|
||||
if _, err := svc.GetLogStream(stream.ID); err != nil {
|
||||
t.Fatalf("get log stream: %v", err)
|
||||
}
|
||||
streams, err := svc.ListLogStreams(domain.LogStreamFilter{ServerInstanceID: instance.ID})
|
||||
if err != nil || len(streams) != 1 {
|
||||
streams, err = svc.ListLogStreams(domain.LogStreamFilter{ServerInstanceID: instance.ID})
|
||||
if err != nil || len(streams) != 3 {
|
||||
t.Fatalf("list log streams: len=%d err=%v", len(streams), err)
|
||||
}
|
||||
|
||||
@@ -183,6 +187,109 @@ func TestCoreServiceCreateListGetWorkflows(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestCoreServiceCreateRemoteProgramJobCreatesManagementLogStreams(t *testing.T) {
|
||||
svc := newTestCoreService()
|
||||
plugin, endpoint := createPluginAndRunEndpoint(t, svc)
|
||||
plugin.RequiredRunCapabilities = append(plugin.RequiredRunCapabilities, domain.JobCapabilityRemoteRunProgram)
|
||||
if err := svc.store.GamePlugins().Update(plugin); err != nil {
|
||||
t.Fatalf("update plugin capabilities: %v", err)
|
||||
}
|
||||
endpoint.Capabilities = append(endpoint.Capabilities, domain.JobCapabilityRemoteRunProgram)
|
||||
if err := svc.store.RunEndpoints().Update(endpoint); err != nil {
|
||||
t.Fatalf("update endpoint capabilities: %v", err)
|
||||
}
|
||||
instance, err := svc.CreateServerInstance(domain.ServerInstance{
|
||||
ID: "server-terminal",
|
||||
PluginID: plugin.ID,
|
||||
RunEndpointID: endpoint.ID,
|
||||
Name: "SCUM Terminal",
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("create server instance: %v", err)
|
||||
}
|
||||
job, err := svc.CreateJob(domain.Job{
|
||||
ID: "job-terminal",
|
||||
ServerInstanceID: instance.ID,
|
||||
RunEndpointID: endpoint.ID,
|
||||
Capability: domain.JobCapabilityRemoteRunProgram,
|
||||
TargetKey: "protected-program",
|
||||
InputRef: "input://protected-program/job-terminal",
|
||||
IdempotencyKey: "idem-terminal",
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("create remote program job: %v", err)
|
||||
}
|
||||
streams, err := svc.ListLogStreams(domain.LogStreamFilter{ServerInstanceID: instance.ID})
|
||||
if err != nil {
|
||||
t.Fatalf("list log streams: %v", err)
|
||||
}
|
||||
if len(streams) != 4 {
|
||||
t.Fatalf("expected stdout/stderr plus management program streams, got %+v", streams)
|
||||
}
|
||||
want := map[string]domain.LogStreamSource{
|
||||
"stdout": domain.LogStreamSourceProcess,
|
||||
"stderr": domain.LogStreamSourceProcess,
|
||||
"management-program.stdout": domain.LogStreamSourceManagementProgram,
|
||||
"management-program.stderr": domain.LogStreamSourceManagementProgram,
|
||||
}
|
||||
for _, stream := range streams {
|
||||
source, ok := want[stream.StreamKey]
|
||||
if !ok {
|
||||
t.Fatalf("unexpected stream key: %+v", stream)
|
||||
}
|
||||
if stream.Source != source || stream.ID != jobLogStreamID(job.ID, stream.StreamKey) {
|
||||
t.Fatalf("unexpected stream metadata: %+v", stream)
|
||||
}
|
||||
delete(want, stream.StreamKey)
|
||||
}
|
||||
if len(want) != 0 {
|
||||
t.Fatalf("missing streams: %+v", want)
|
||||
}
|
||||
}
|
||||
|
||||
func TestCoreServiceStartupRecoversLegacyJobLogStreams(t *testing.T) {
|
||||
store := repo.NewMemoryStore()
|
||||
seed := newCoreService(store, func() time.Time { return fixedTime })
|
||||
plugin, endpoint := createPluginAndRunEndpoint(t, seed)
|
||||
plugin.RequiredRunCapabilities = append(plugin.RequiredRunCapabilities, domain.JobCapabilityRemoteRunProgram)
|
||||
if err := seed.store.GamePlugins().Update(plugin); err != nil {
|
||||
t.Fatalf("update plugin capabilities: %v", err)
|
||||
}
|
||||
endpoint.Capabilities = append(endpoint.Capabilities, domain.JobCapabilityRemoteRunProgram)
|
||||
if err := seed.store.RunEndpoints().Update(endpoint); err != nil {
|
||||
t.Fatalf("update endpoint capabilities: %v", err)
|
||||
}
|
||||
if _, err := seed.CreateServerInstance(domain.ServerInstance{ID: "legacy-terminal-server", PluginID: plugin.ID, RunEndpointID: endpoint.ID, Name: "Legacy Terminal"}); err != nil {
|
||||
t.Fatalf("create server instance: %v", err)
|
||||
}
|
||||
if err := store.Jobs().Create(domain.Job{
|
||||
ID: "legacy-terminal-job",
|
||||
ServerInstanceID: "legacy-terminal-server",
|
||||
RunEndpointID: endpoint.ID,
|
||||
Capability: domain.JobCapabilityRemoteRunProgram,
|
||||
TargetKey: "protected-program",
|
||||
InputRef: "input://protected-program/legacy-terminal-job",
|
||||
IdempotencyKey: "legacy-terminal",
|
||||
State: domain.JobStateQueued,
|
||||
CreatedAt: fixedTime,
|
||||
UpdatedAt: fixedTime,
|
||||
}); err != nil {
|
||||
t.Fatalf("seed legacy job: %v", err)
|
||||
}
|
||||
before, err := seed.ListLogStreams(domain.LogStreamFilter{ServerInstanceID: "legacy-terminal-server"})
|
||||
if err != nil || len(before) != 0 {
|
||||
t.Fatalf("expected no seeded streams before recovery, len=%d err=%v streams=%+v", len(before), err, before)
|
||||
}
|
||||
recovered, err := NewCoreServiceWithDurableStores(store, NewMemoryLogBodyStore(), NewMemoryArtifactBodyStore())
|
||||
if err != nil {
|
||||
t.Fatalf("recover durable service: %v", err)
|
||||
}
|
||||
after, err := recovered.ListLogStreams(domain.LogStreamFilter{ServerInstanceID: "legacy-terminal-server"})
|
||||
if err != nil || len(after) != 4 {
|
||||
t.Fatalf("expected recovered job log streams, len=%d err=%v streams=%+v", len(after), err, after)
|
||||
}
|
||||
}
|
||||
|
||||
func TestCoreServiceRejectsInvalidServerDependencies(t *testing.T) {
|
||||
svc := newTestCoreService()
|
||||
plugin, endpoint := createPluginAndRunEndpoint(t, svc)
|
||||
|
||||
@@ -68,7 +68,7 @@ func ValidateLogBatchIngest(batch domain.LogBatchIngest) error {
|
||||
computed, err := LogEntriesChecksum(batch.Entries)
|
||||
if err != nil {
|
||||
violations = append(violations, "checksum cannot be computed")
|
||||
} else if batch.Checksum != computed {
|
||||
} else if batch.Checksum != computed && !logLineChecksumMatches(batch) {
|
||||
violations = append(violations, "checksum does not match entries")
|
||||
}
|
||||
}
|
||||
@@ -107,6 +107,18 @@ func LogEntriesChecksum(entries []domain.LogEntry) (string, error) {
|
||||
return "sha256:" + hex.EncodeToString(sum[:]), nil
|
||||
}
|
||||
|
||||
func logLineChecksumMatches(batch domain.LogBatchIngest) bool {
|
||||
if len(batch.Entries) != 1 {
|
||||
return false
|
||||
}
|
||||
return batch.Checksum == LogLineChecksum(batch.Entries[0].Line)
|
||||
}
|
||||
|
||||
func LogLineChecksum(value string) string {
|
||||
sum := sha256.Sum256([]byte(value))
|
||||
return "sha256:" + hex.EncodeToString(sum[:])
|
||||
}
|
||||
|
||||
type logEntryChecksumBody struct {
|
||||
Seq uint64 `json:"seq"`
|
||||
Timestamp string `json:"timestamp"`
|
||||
|
||||
@@ -2355,7 +2355,7 @@ func validArtifactState(state domain.ArtifactState) bool {
|
||||
|
||||
func validLogStreamSource(source domain.LogStreamSource) bool {
|
||||
switch source {
|
||||
case domain.LogStreamSourceProcess, domain.LogStreamSourceFile, domain.LogStreamSourcePlugin:
|
||||
case domain.LogStreamSourceProcess, domain.LogStreamSourceFile, domain.LogStreamSourcePlugin, domain.LogStreamSourceManagementProgram:
|
||||
return true
|
||||
default:
|
||||
return strings.TrimSpace(string(source)) != ""
|
||||
|
||||
Reference in New Issue
Block a user