Stream live server logs over SSE

This commit is contained in:
npc0-hue
2026-08-03 22:28:54 +08:00
parent 5d4fca14f9
commit 7eac1926dd
48 changed files with 1526 additions and 263 deletions
+119
View File
@@ -39,6 +39,38 @@ func TestCoreServiceIngestsLogBatchAndQueriesCursor(t *testing.T) {
}
}
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.Stream.ID != "log-1" || event.Entry.Seq != 1 || event.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 TestCoreServiceLogBatchDuplicateAck(t *testing.T) {
svc, sessionToken := newRegisteredLogIngestService(t)
createLogStreamFixture(t, svc)
@@ -114,6 +146,93 @@ func TestCoreServiceAcceptsAutoCreatedRunJobLogStreams(t *testing.T) {
}
}
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 TestCoreServiceRejectsOutOfOrderAndConflictingLogBatches(t *testing.T) {
svc, sessionToken := newRegisteredLogIngestService(t)
createLogStreamFixture(t, svc)