603 lines
23 KiB
Go
603 lines
23 KiB
Go
package service
|
|
|
|
import (
|
|
"path/filepath"
|
|
"strings"
|
|
"sync"
|
|
"testing"
|
|
"time"
|
|
|
|
"browser.local/platform/domain"
|
|
"browser.local/platform/validator"
|
|
)
|
|
|
|
func TestCoreServiceIngestsLogBatchAndQueriesCursor(t *testing.T) {
|
|
svc, sessionToken := newRegisteredLogIngestService(t)
|
|
createLogStreamFixture(t, svc)
|
|
batch := validLogBatch(t, sessionToken, 1, 2)
|
|
|
|
ack, err := svc.IngestLogBatch(batch)
|
|
if err != nil {
|
|
t.Fatalf("ingest log batch: %v", err)
|
|
}
|
|
if !ack.Accepted || ack.AcceptedFrom != 1 || ack.AcceptedTo != 2 || ack.LatestSeq != 2 {
|
|
t.Fatalf("unexpected ack: %+v", ack)
|
|
}
|
|
stream, err := svc.GetLogStream("log-1")
|
|
if err != nil {
|
|
t.Fatalf("get log stream: %v", err)
|
|
}
|
|
if stream.LatestSeq != 2 {
|
|
t.Fatalf("expected latest seq 2, got %+v", stream)
|
|
}
|
|
|
|
query, err := svc.QueryLogStream(domain.LogStreamCursorQuery{LogStreamID: "log-1", AfterSeq: 1, Limit: 10})
|
|
if err != nil {
|
|
t.Fatalf("query log stream: %v", err)
|
|
}
|
|
if len(query.Entries) != 1 || query.Entries[0].Seq != 2 || query.NextSeq != 2 || query.LatestSeq != 2 {
|
|
t.Fatalf("unexpected query result: %+v", query)
|
|
}
|
|
}
|
|
|
|
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)
|
|
subscription, err := svc.SubscribeLogEvents("server-1")
|
|
if err != nil {
|
|
t.Fatalf("subscribe log events: %v", err)
|
|
}
|
|
defer subscription.Close()
|
|
|
|
batch := validLogBatch(t, sessionToken, 1, 1)
|
|
if _, err := svc.IngestLogBatch(batch); err != nil {
|
|
t.Fatalf("ingest log batch: %v", err)
|
|
}
|
|
select {
|
|
case event := <-subscription.Events:
|
|
if event.Kind != LogEventSubscriptionEventLog || event.LogEvent.Stream.ID != "log-1" || event.LogEvent.Entry.Seq != 1 || event.LogEvent.LatestSeq != 1 {
|
|
t.Fatalf("unexpected log event: %+v", event)
|
|
}
|
|
case <-time.After(time.Second):
|
|
t.Fatal("expected log event after accepted batch")
|
|
}
|
|
|
|
if _, err := svc.IngestLogBatch(batch); err != nil {
|
|
t.Fatalf("ingest duplicate batch: %v", err)
|
|
}
|
|
select {
|
|
case event := <-subscription.Events:
|
|
t.Fatalf("duplicate batch should not publish a second event: %+v", event)
|
|
default:
|
|
}
|
|
}
|
|
|
|
func TestCoreServicePersistsAndEnforcesImmutableProcessLogSessionMetadata(t *testing.T) {
|
|
svc, sessionToken := newRegisteredLogIngestService(t)
|
|
startedAt := time.Date(2026, 7, 3, 12, 30, 0, 0, time.UTC)
|
|
streamID := runSessionLogStreamID("run-local", "server-1", "session-a", "stdout")
|
|
batch := validLogBatch(t, sessionToken, 1, 1)
|
|
batch.LogStreamID = streamID
|
|
batch.LogSessionID = "session-a"
|
|
batch.SessionStartedAt = startedAt
|
|
if _, err := svc.IngestLogBatch(batch); err != nil {
|
|
t.Fatalf("ingest session-scoped process batch: %v", err)
|
|
}
|
|
stream, err := svc.GetLogStream(streamID)
|
|
if err != nil || stream.LogSessionID != "session-a" || !stream.SessionStartedAt.Equal(startedAt) {
|
|
t.Fatalf("session metadata was not persisted: stream=%+v err=%v", stream, err)
|
|
}
|
|
|
|
conflict := validLogBatch(t, sessionToken, 2, 2)
|
|
conflict.LogStreamID = streamID
|
|
conflict.LogSessionID = "session-a"
|
|
conflict.SessionStartedAt = startedAt.Add(time.Second)
|
|
if _, err := svc.IngestLogBatch(conflict); err == nil || !strings.Contains(err.Error(), "metadata must match") {
|
|
t.Fatalf("expected immutable stream metadata rejection, got %v", err)
|
|
}
|
|
|
|
legacy := createLogStreamFixture(t, svc)
|
|
legacyBatch := validLogBatch(t, sessionToken, 1, 1)
|
|
legacyBatch.LogStreamID = legacy.ID
|
|
legacyBatch.LogSessionID = "session-a"
|
|
legacyBatch.SessionStartedAt = startedAt
|
|
if _, err := svc.IngestLogBatch(legacyBatch); err == nil || !strings.Contains(err.Error(), "metadata must match") {
|
|
t.Fatalf("expected legacy stream to reject attached session metadata, got %v", err)
|
|
}
|
|
}
|
|
|
|
func TestCoreServiceAcceptsLegacyProcessStreamStartWithoutSessionID(t *testing.T) {
|
|
svc, sessionToken := newRegisteredLogIngestService(t)
|
|
startedAt := time.Date(2026, 7, 3, 12, 30, 0, 0, time.UTC)
|
|
streamID := "run.run-local.server-1.stdout"
|
|
if _, err := svc.CreateLogStream(domain.LogStream{
|
|
ID: streamID,
|
|
ServerInstanceID: "server-1",
|
|
Source: domain.LogStreamSourceProcess,
|
|
StreamKey: "stdout",
|
|
SessionStartedAt: startedAt,
|
|
StorageBackend: domain.LogStorageBackendLocalSegments,
|
|
RetentionPolicy: "default",
|
|
}); err != nil {
|
|
t.Fatalf("create legacy process stream: %v", err)
|
|
}
|
|
|
|
batch := validLogBatch(t, sessionToken, 1, 1)
|
|
batch.LogStreamID = streamID
|
|
batch.SessionStartedAt = startedAt
|
|
if _, err := svc.IngestLogBatch(batch); err != nil {
|
|
t.Fatalf("ingest legacy process batch: %v", err)
|
|
}
|
|
}
|
|
|
|
func TestCoreServiceSerializesConsistentSessionStreamCreation(t *testing.T) {
|
|
svc, _ := newRegisteredLogIngestService(t)
|
|
startedAt := time.Date(2026, 7, 3, 12, 30, 0, 0, time.UTC)
|
|
streams := []domain.LogStream{
|
|
{ID: "session-stream-stdout", ServerInstanceID: "server-1", Source: domain.LogStreamSourceProcess, StreamKey: "stdout", LogSessionID: "session-a", SessionStartedAt: startedAt, StorageBackend: domain.LogStorageBackendLocalSegments, RetentionPolicy: "default"},
|
|
{ID: "session-stream-stderr", ServerInstanceID: "server-1", Source: domain.LogStreamSourceProcess, StreamKey: "stderr", LogSessionID: "session-a", SessionStartedAt: startedAt, StorageBackend: domain.LogStorageBackendLocalSegments, RetentionPolicy: "default"},
|
|
}
|
|
errorsByStream := make(chan error, len(streams))
|
|
var wait sync.WaitGroup
|
|
for _, stream := range streams {
|
|
stream := stream
|
|
wait.Add(1)
|
|
go func() {
|
|
defer wait.Done()
|
|
_, err := svc.CreateLogStream(stream)
|
|
errorsByStream <- err
|
|
}()
|
|
}
|
|
wait.Wait()
|
|
close(errorsByStream)
|
|
for err := range errorsByStream {
|
|
if err != nil {
|
|
t.Fatalf("create consistent session stream: %v", err)
|
|
}
|
|
}
|
|
|
|
_, err := svc.CreateLogStream(domain.LogStream{ID: "session-stream-conflict", ServerInstanceID: "server-1", Source: domain.LogStreamSourceProcess, StreamKey: "console", LogSessionID: "session-a", SessionStartedAt: startedAt.Add(time.Second), StorageBackend: domain.LogStorageBackendLocalSegments, RetentionPolicy: "default"})
|
|
if err == nil || !strings.Contains(err.Error(), "conflicts") {
|
|
t.Fatalf("expected conflicting session timestamp rejection, got %v", err)
|
|
}
|
|
_, err = svc.CreateLogStream(domain.LogStream{ID: "session-file-tail", ServerInstanceID: "server-1", Source: domain.LogStreamSourceFile, StreamKey: "file", LogSessionID: "session-file", SessionStartedAt: startedAt, StorageBackend: domain.LogStorageBackendLocalSegments, RetentionPolicy: "default"})
|
|
if err == nil || !strings.Contains(err.Error(), "only valid for process") {
|
|
t.Fatalf("expected file-tail session metadata rejection, got %v", err)
|
|
}
|
|
}
|
|
|
|
func TestCoreServiceLogBatchDuplicateAck(t *testing.T) {
|
|
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)
|
|
}
|
|
ack, err := svc.IngestLogBatch(batch)
|
|
if err != nil {
|
|
t.Fatalf("ingest duplicate batch: %v", err)
|
|
}
|
|
if !ack.Duplicate || ack.LatestSeq != 2 {
|
|
t.Fatalf("expected duplicate ack, got %+v", ack)
|
|
}
|
|
}
|
|
|
|
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 TestCoreServiceAcceptsPluginDeclaredProcessLogStreams(t *testing.T) {
|
|
svc, sessionToken := newRegisteredLogIngestService(t)
|
|
job, err := svc.CreateJob(domain.Job{
|
|
ID: "job-declared-process-logs",
|
|
ServerInstanceID: "server-1",
|
|
RunEndpointID: "run-local",
|
|
Capability: domain.LifecycleCapabilityStart,
|
|
IdempotencyKey: "job-declared-process-logs",
|
|
ExecutionInput: domain.JobExecutionInput{WorkspaceScope: "run-local", LifecycleOperation: "start", LogSources: []domain.RuntimeLogSource{
|
|
{Key: "console-out", Kind: "process.stdout", StreamKey: "scum.console.stdout", CursorKind: "sequence", RetentionDays: 30},
|
|
{Key: "console-err", Kind: "process.stderr", StreamKey: "scum.console.stderr", CursorKind: "sequence", RetentionDays: 30},
|
|
}},
|
|
})
|
|
if err != nil {
|
|
t.Fatalf("create job: %v", err)
|
|
}
|
|
streamID := jobLogStreamID(job.ID, "scum.console.stdout")
|
|
stream, err := svc.GetLogStream(streamID)
|
|
if err != nil {
|
|
t.Fatalf("get declared process log stream: %v", err)
|
|
}
|
|
if stream.StreamKey != "scum.console.stdout" || stream.Source != domain.LogStreamSourceProcess {
|
|
t.Fatalf("unexpected declared stream metadata: %+v", stream)
|
|
}
|
|
entry := domain.LogEntry{Seq: 1, Timestamp: time.Date(2026, 7, 3, 12, 0, 1, 0, time.UTC), Level: "info", Line: "LogStreaming: Display: server ready"}
|
|
ack, err := svc.IngestLogBatch(domain.LogBatchIngest{
|
|
RunEndpointID: "run-local",
|
|
SessionToken: sessionToken,
|
|
LogStreamID: streamID,
|
|
ServerInstanceID: "server-1",
|
|
StreamKey: "scum.console.stdout",
|
|
Source: domain.LogStreamSourceProcess,
|
|
FirstSeq: entry.Seq,
|
|
LastSeq: entry.Seq,
|
|
Compression: "none",
|
|
Checksum: validator.LogLineChecksum(entry.Line),
|
|
Entries: []domain.LogEntry{entry},
|
|
})
|
|
if err != nil {
|
|
t.Fatalf("ingest declared process log batch: %v", err)
|
|
}
|
|
if !ack.Accepted || ack.LatestSeq != entry.Seq {
|
|
t.Fatalf("unexpected declared stream ack: %+v", ack)
|
|
}
|
|
}
|
|
|
|
func TestCoreServiceRepairsMissingDeclaredProcessLogStreamOnIngest(t *testing.T) {
|
|
svc, sessionToken := newRegisteredLogIngestService(t)
|
|
job := domain.Job{
|
|
ID: "job-repaired-process-logs",
|
|
ServerInstanceID: "server-1",
|
|
RunEndpointID: "run-local",
|
|
Capability: domain.LifecycleCapabilityStart,
|
|
IdempotencyKey: "job-repaired-process-logs",
|
|
ExecutionInput: domain.JobExecutionInput{WorkspaceScope: "run-local", LifecycleOperation: "start", LogSources: []domain.RuntimeLogSource{
|
|
{Key: "console-out", Kind: "process.stdout", StreamKey: "scum.console.stdout", CursorKind: "sequence", RetentionDays: 30},
|
|
}},
|
|
}
|
|
if err := svc.store.Jobs().Create(job); err != nil {
|
|
t.Fatalf("seed legacy job without streams: %v", err)
|
|
}
|
|
streamID := jobLogStreamID(job.ID, "scum.console.stdout")
|
|
if _, err := svc.GetLogStream(streamID); err == nil {
|
|
t.Fatal("expected declared stream to be missing before ingest repair")
|
|
}
|
|
entry := domain.LogEntry{Seq: 1, Timestamp: time.Date(2026, 7, 3, 12, 0, 1, 0, time.UTC), Level: "info", Line: "LogStreaming: Display: recovered from spool"}
|
|
ack, err := svc.IngestLogBatch(domain.LogBatchIngest{
|
|
RunEndpointID: "run-local",
|
|
SessionToken: sessionToken,
|
|
LogStreamID: streamID,
|
|
ServerInstanceID: "server-1",
|
|
StreamKey: "scum.console.stdout",
|
|
Source: domain.LogStreamSourceProcess,
|
|
FirstSeq: entry.Seq,
|
|
LastSeq: entry.Seq,
|
|
Compression: "none",
|
|
Checksum: validator.LogLineChecksum(entry.Line),
|
|
Entries: []domain.LogEntry{entry},
|
|
})
|
|
if err != nil {
|
|
t.Fatalf("ingest repaired declared process log batch: %v", err)
|
|
}
|
|
if !ack.Accepted || ack.LatestSeq != entry.Seq {
|
|
t.Fatalf("unexpected repaired stream ack: %+v", ack)
|
|
}
|
|
}
|
|
|
|
func TestCoreServiceAcceptsAutonomousRunLogStreamWithoutPlatformJob(t *testing.T) {
|
|
svc, sessionToken := newRegisteredLogIngestService(t)
|
|
entry := domain.LogEntry{Seq: 1, Timestamp: time.Date(2026, 7, 3, 12, 0, 1, 0, time.UTC), Level: "info", Line: "autonomous bootstrap output"}
|
|
streamID := runLogStreamID("run-local", "server-1", "scum.console.stdout")
|
|
ack, err := svc.IngestLogBatch(domain.LogBatchIngest{
|
|
RunEndpointID: "run-local",
|
|
SessionToken: sessionToken,
|
|
LogStreamID: streamID,
|
|
ServerInstanceID: "server-1",
|
|
StreamKey: "scum.console.stdout",
|
|
Source: domain.LogStreamSourceProcess,
|
|
FirstSeq: entry.Seq,
|
|
LastSeq: entry.Seq,
|
|
Compression: "none",
|
|
Checksum: validator.LogLineChecksum(entry.Line),
|
|
Entries: []domain.LogEntry{entry},
|
|
})
|
|
if err != nil {
|
|
t.Fatalf("ingest autonomous run log batch: %v", err)
|
|
}
|
|
if !ack.Accepted || ack.LogStreamID != streamID || ack.LatestSeq != entry.Seq {
|
|
t.Fatalf("unexpected autonomous stream ack: %+v", ack)
|
|
}
|
|
stream, err := svc.GetLogStream(streamID)
|
|
if err != nil {
|
|
t.Fatalf("get autonomous stream: %v", err)
|
|
}
|
|
if stream.ServerInstanceID != "server-1" || stream.StreamKey != "scum.console.stdout" || stream.Source != domain.LogStreamSourceProcess {
|
|
t.Fatalf("unexpected autonomous stream metadata: %+v", stream)
|
|
}
|
|
}
|
|
|
|
func TestCoreServiceAcceptsAutonomousRunFileTailLogStreamWithoutPlatformJob(t *testing.T) {
|
|
svc, sessionToken := newRegisteredLogIngestService(t)
|
|
entry := domain.LogEntry{Seq: 1, Timestamp: time.Date(2026, 7, 3, 12, 0, 1, 0, time.UTC), Level: "info", Line: "LogSCUM: Display: player joined"}
|
|
streamID := runLogStreamID("run-local", "server-1", "scum.server")
|
|
ack, err := svc.IngestLogBatch(domain.LogBatchIngest{
|
|
RunEndpointID: "run-local",
|
|
SessionToken: sessionToken,
|
|
LogStreamID: streamID,
|
|
ServerInstanceID: "server-1",
|
|
StreamKey: "scum.server",
|
|
Source: domain.LogStreamSourceFile,
|
|
FirstSeq: entry.Seq,
|
|
LastSeq: entry.Seq,
|
|
Compression: "none",
|
|
Checksum: validator.LogLineChecksum(entry.Line),
|
|
Entries: []domain.LogEntry{entry},
|
|
})
|
|
if err != nil {
|
|
t.Fatalf("ingest autonomous run file-tail log batch: %v", err)
|
|
}
|
|
if !ack.Accepted || ack.LogStreamID != streamID || ack.LatestSeq != entry.Seq {
|
|
t.Fatalf("unexpected autonomous file-tail stream ack: %+v", ack)
|
|
}
|
|
stream, err := svc.GetLogStream(streamID)
|
|
if err != nil {
|
|
t.Fatalf("get autonomous file-tail stream: %v", err)
|
|
}
|
|
if stream.ServerInstanceID != "server-1" || stream.StreamKey != "scum.server" || stream.Source != domain.LogStreamSourceFile {
|
|
t.Fatalf("unexpected autonomous file-tail stream metadata: %+v", stream)
|
|
}
|
|
query, err := svc.QueryLogStream(domain.LogStreamCursorQuery{LogStreamID: streamID, AfterSeq: 0, Limit: 10})
|
|
if err != nil || len(query.Entries) != 1 || query.Entries[0].Line != entry.Line {
|
|
t.Fatalf("unexpected autonomous file-tail query: %+v err=%v", query, err)
|
|
}
|
|
}
|
|
|
|
func TestCoreServiceAcceptsLegacyAutonomousJobLogStreamWithoutPlatformJob(t *testing.T) {
|
|
svc, sessionToken := newRegisteredLogIngestService(t)
|
|
entry := domain.LogEntry{Seq: 1, Timestamp: time.Date(2026, 7, 3, 12, 0, 1, 0, time.UTC), Level: "info", Line: "legacy autonomous bootstrap output"}
|
|
streamID := jobLogStreamID("autonomous-bootstrap-start", "scum.console.stdout")
|
|
ack, err := svc.IngestLogBatch(domain.LogBatchIngest{
|
|
RunEndpointID: "run-local",
|
|
SessionToken: sessionToken,
|
|
LogStreamID: streamID,
|
|
ServerInstanceID: "server-1",
|
|
StreamKey: "scum.console.stdout",
|
|
Source: domain.LogStreamSourceProcess,
|
|
FirstSeq: entry.Seq,
|
|
LastSeq: entry.Seq,
|
|
Compression: "none",
|
|
Checksum: validator.LogLineChecksum(entry.Line),
|
|
Entries: []domain.LogEntry{entry},
|
|
})
|
|
if err != nil {
|
|
t.Fatalf("ingest legacy autonomous log batch: %v", err)
|
|
}
|
|
if !ack.Accepted || ack.LogStreamID != streamID || ack.LatestSeq != entry.Seq {
|
|
t.Fatalf("unexpected legacy autonomous stream ack: %+v", ack)
|
|
}
|
|
}
|
|
|
|
func TestCoreServiceRejectsSessionMetadataOnLegacyAutonomousStreamID(t *testing.T) {
|
|
svc, sessionToken := newRegisteredLogIngestService(t)
|
|
batch := validLogBatch(t, sessionToken, 1, 1)
|
|
batch.LogStreamID = jobLogStreamID("autonomous-bootstrap-start", "stdout")
|
|
batch.LogSessionID = "session-a"
|
|
batch.SessionStartedAt = time.Date(2026, 7, 3, 12, 30, 0, 0, time.UTC)
|
|
if _, err := svc.IngestLogBatch(batch); err == nil {
|
|
t.Fatal("expected session-scoped batch with legacy autonomous stream ID to be rejected")
|
|
}
|
|
}
|
|
|
|
func TestCoreServiceRejectsOutOfOrderAndConflictingLogBatches(t *testing.T) {
|
|
svc, sessionToken := newRegisteredLogIngestService(t)
|
|
createLogStreamFixture(t, svc)
|
|
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)
|
|
}
|
|
conflict := batch
|
|
conflict.Entries[0].Line = "changed"
|
|
conflict.Checksum = checksumForEntries(t, conflict.Entries)
|
|
_, err = svc.IngestLogBatch(conflict)
|
|
if err == nil || !strings.Contains(err.Error(), "conflicts") {
|
|
t.Fatalf("expected conflicting duplicate rejection, got %v", err)
|
|
}
|
|
}
|
|
|
|
func TestCoreServiceRejectsMissingLogStream(t *testing.T) {
|
|
svc, sessionToken := newRegisteredLogIngestService(t)
|
|
_, err := svc.IngestLogBatch(validLogBatch(t, sessionToken, 1, 1))
|
|
if err == nil {
|
|
t.Fatal("expected missing stream error")
|
|
}
|
|
}
|
|
|
|
func TestFileLogBodyStoreReloadsBatchesAndCursorEntries(t *testing.T) {
|
|
rootDir := filepath.Join(t.TempDir(), "logs")
|
|
store, err := NewFileLogBodyStore(rootDir)
|
|
if err != nil {
|
|
t.Fatalf("create file log store: %v", err)
|
|
}
|
|
entries := []domain.LogEntry{
|
|
{Seq: 1, Timestamp: time.Date(2026, 7, 3, 12, 0, 1, 0, time.UTC), Level: "info", Line: "one"},
|
|
{Seq: 2, Timestamp: time.Date(2026, 7, 3, 12, 0, 2, 0, time.UTC), Level: "warn", Line: "two"},
|
|
}
|
|
record := domain.LogBatchRecord{
|
|
Checksum: checksumForEntries(t, entries),
|
|
FirstSeq: 1,
|
|
LastSeq: 2,
|
|
Entries: entries,
|
|
}
|
|
if err := store.AppendBatch("log-1", record); err != nil {
|
|
t.Fatalf("append batch: %v", err)
|
|
}
|
|
|
|
reloaded, err := NewFileLogBodyStore(rootDir)
|
|
if err != nil {
|
|
t.Fatalf("reload file log store: %v", err)
|
|
}
|
|
got, exists, err := reloaded.GetBatch("log-1", 1)
|
|
if err != nil {
|
|
t.Fatalf("get reloaded batch: %v", err)
|
|
}
|
|
if !exists || got.Checksum != record.Checksum || got.LastSeq != 2 {
|
|
t.Fatalf("unexpected reloaded batch: exists=%v record=%+v", exists, got)
|
|
}
|
|
selected, nextSeq, err := reloaded.Query("log-1", 1, 10)
|
|
if err != nil {
|
|
t.Fatalf("query reloaded entries: %v", err)
|
|
}
|
|
if len(selected) != 1 || selected[0].Seq != 2 || selected[0].Line != "two" || nextSeq != 2 {
|
|
t.Fatalf("unexpected reloaded query: entries=%+v next=%d", selected, nextSeq)
|
|
}
|
|
}
|
|
|
|
func newRegisteredLogIngestService(t *testing.T) (*CoreService, string) {
|
|
t.Helper()
|
|
svc := newTestCoreService()
|
|
plugin, endpoint := createPluginAndRunEndpoint(t, svc)
|
|
if _, err := svc.CreateServerInstance(domain.ServerInstance{
|
|
ID: "server-1",
|
|
PluginID: plugin.ID,
|
|
RunEndpointID: endpoint.ID,
|
|
Name: "SCUM #1",
|
|
}); err != nil {
|
|
t.Fatalf("create server instance: %v", err)
|
|
}
|
|
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)
|
|
}
|
|
return svc, hello.SessionToken
|
|
}
|
|
|
|
func createLogStreamFixture(t *testing.T, svc *CoreService) domain.LogStream {
|
|
t.Helper()
|
|
stream, err := svc.CreateLogStream(domain.LogStream{
|
|
ID: "log-1",
|
|
ServerInstanceID: "server-1",
|
|
Source: domain.LogStreamSourceProcess,
|
|
StreamKey: "stdout",
|
|
StorageBackend: domain.LogStorageBackendLocalSegments,
|
|
RetentionPolicy: "default",
|
|
})
|
|
if err != nil {
|
|
t.Fatalf("create log stream: %v", err)
|
|
}
|
|
return stream
|
|
}
|
|
|
|
func validLogBatch(t *testing.T, sessionToken string, firstSeq uint64, lastSeq uint64) domain.LogBatchIngest {
|
|
t.Helper()
|
|
entries := make([]domain.LogEntry, 0, lastSeq-firstSeq+1)
|
|
for seq := firstSeq; seq <= lastSeq; seq++ {
|
|
entries = append(entries, domain.LogEntry{
|
|
Seq: seq,
|
|
Timestamp: time.Date(2026, 7, 3, 12, 0, int(seq), 0, time.UTC),
|
|
Level: "info",
|
|
Line: "line",
|
|
})
|
|
}
|
|
return domain.LogBatchIngest{
|
|
RunEndpointID: "run-local",
|
|
SessionToken: sessionToken,
|
|
LogStreamID: "log-1",
|
|
ServerInstanceID: "server-1",
|
|
StreamKey: "stdout",
|
|
Source: domain.LogStreamSourceProcess,
|
|
FirstSeq: firstSeq,
|
|
LastSeq: lastSeq,
|
|
Compression: "none",
|
|
Checksum: checksumForEntries(t, entries),
|
|
Entries: entries,
|
|
}
|
|
}
|
|
|
|
func checksumForEntries(t *testing.T, entries []domain.LogEntry) string {
|
|
t.Helper()
|
|
checksum, err := validator.LogEntriesChecksum(entries)
|
|
if err != nil {
|
|
t.Fatalf("checksum entries: %v", err)
|
|
}
|
|
return checksum
|
|
}
|