diff --git a/docs/01-basic/report.md b/docs/01-basic/report.md new file mode 100644 index 0000000..9d35a21 --- /dev/null +++ b/docs/01-basic/report.md @@ -0,0 +1,86 @@ +# Report for Basic Functions + +## Q1.1 + +### 按逗号分割与字段含义 + +代码框架并不是手动调用 `String.Split` 来按逗号分割日志,而是使用了 **CsvHelper** 库。在 `LogFileParser.Parse` 方法中: + +```csharp +using var csv = new CsvReader(logFile, config); +csv.Context.RegisterClassMap(); +foreach (var logRecord in csv.GetRecords()) +{ + yield return LineParser.ParseLine(logRecord); +} +``` + +- `new CsvReader(logFile, config)` 创建了 CSV 读取器,`csv.GetRecords()` 负责按逗号把每一行拆分成若干字段。 +- 每一行第几个字段代表什么含义,由 `LogRecordMap`(继承 `CsvHelper.Configuration.ClassMap`)指定: + + ```csharp + Map(m => m.LineNo).Index(0); // 第 0 列是行号 + Map(m => m.Timestamp).Index(1); // 第 1 列是时间戳 + Map(m => m.PodName).Index(2); // 第 2 列是容器名 + Map(m => m.Message).Index(3); // 第 3 列是 JSON 消息 + ``` + + 并通过 `csv.Context.RegisterClassMap()` 注册生效。 + +### 判断日志种类 + +在 `LineParser.ParseLine` 方法中,通过以下语句判断这一行日志的种类: + +```csharp +using var doc = JsonDocument.Parse(logRecord.Message); +var root = doc.RootElement; +if (root.TryGetProperty("event", out var eventElement)) +{ + return eventElement.GetString() switch + { + "call" => LineParser.CreateCall(logRecord), + "request" => LineParser.CreateRequest(logRecord), + "internal" => LineParser.CreateInternal(logRecord), + _ => throw new FormatException(...) + }; +} +``` + +即先用 `JsonDocument.Parse` 解析 `message`,再用 `root.TryGetProperty("event", ...)` 取出 `event` 字段,最后用 `switch` 表达式根据 `event` 的值(`"call"` / `"request"` / `"internal"`)分发到对应的创建方法。 + +### 解析 JSON 所用的库方法 + +确定日志种类后,调用的是 `System.Text.Json` 中的 `JsonSerializer.Deserialize(json, options)`(例如 `JsonSerializer.Deserialize(logRecord.Message, options)`)。 + +**防止字段缺失:** 每个 `Message` record 的字段都标注了 `[property: JsonRequired]` 特性(例如 `[property: JsonRequired] string RequestId`)。当 JSON 中缺少被标记的字段时,`JsonSerializer.Deserialize` 会抛出 `JsonException`;此外还通过 `?? throw new FormatException(...)` 处理反序列化结果整体为 `null` 的情况。 + +**命名法转换:** 通过 `JsonSerializerOptions` 的命名策略完成: + +```csharp +private static JsonSerializerOptions options = new JsonSerializerOptions +{ + PropertyNamingPolicy = JsonNamingPolicy.KebabCaseLower, +}; +``` + +`PropertyNamingPolicy = JsonNamingPolicy.KebabCaseLower` 会让序列化器把大驼峰属性名(如 `RequestId`、`TargetService`、`DurationMs`)自动转换为烤串命名法(`request-id`、`target-service`、`duration-ms`)去匹配 JSON 中的键。 + +## Q1.2 + +以一个 Call 事件的解析结果为例,`Dump` 方法被调用后的方法调用链如下: + ++ `Dictionary KeyValueVisitor.Dump(LogEntry entry)` ++ `TResult CallLogEntry.Accept(ILogEntryVisitor visitor)`(经 `entry.Accept(this)` 多态调用) ++ `Dictionary KeyValueVisitor.Visit(CallLogEntry entry)` + +## Q1.3 + +(本问为个人反思题,请根据你的实际情况选择作答。下方以 Q1.3.b 为例。) + +### Q1.3.b + +本次作业我使用了 AI 辅助完成。我给予 AI 的提示词大致为:「帮我完成 01-basic 要求的所有作业」,并在此之前通过提问明确了当前分支、任务内容与需要改动的文件。 + +与完全依靠传统搜索引擎和自己能力写出的解答相比,AI 的解答好在:能快速梳理出代码框架中需要补全的 `TODO` 位置,并给出与已有 `Call` 实现风格一致的参考代码,节省了大量阅读与试错的时间。 + +但 AI 的解答也存在不足:例如它有时会忽略 `internal` 日志中 `exception` 字段需要按「冒号加空格」拆分成 `ExceptionName` 与 `ExceptionMessage` 这一细节,需要人工结合测试用例(`TestParseInternalLogExampleFailed`)来确认异常格式的处理方式。因此最终代码仍需要人工 review 与验证。 diff --git a/docs/02-multithreading/assets/localcli-full.png b/docs/02-multithreading/assets/localcli-full.png new file mode 100644 index 0000000..8923760 Binary files /dev/null and b/docs/02-multithreading/assets/localcli-full.png differ diff --git a/docs/02-multithreading/assets/localcli-robustness.png b/docs/02-multithreading/assets/localcli-robustness.png new file mode 100644 index 0000000..88a5443 Binary files /dev/null and b/docs/02-multithreading/assets/localcli-robustness.png differ diff --git a/docs/02-multithreading/report.md b/docs/02-multithreading/report.md new file mode 100644 index 0000000..d8636b8 --- /dev/null +++ b/docs/02-multithreading/report.md @@ -0,0 +1,84 @@ +# Report for Multithreading + +## 功能实现简介 + +本节在 `01-basic` 的基础上,实现了目录级别的并行日志分析器,共完成三个部分: + +1. **线程安全队列 `WorkQueue`**(`LogAnalyzer/WorkQueue.cs`):基于 `Queue` + `lock` + `Monitor`(条件变量)实现的、支持"结束放入"操作的无限容量生产者-消费者队列。 +2. **并行日志分析器 `LogFileAnalyzer`**(`LogAnalyzer/LogFileAnalyzer.cs`):扫描指定目录下所有 `.log` 文件,开多个工作线程并行解析,并缓存每个文件的分析结果。 +3. **控制台交互界面 `LocalCli`**(`LocalCli/Program.cs`):提供展示文件、分析指定文件、分析全部文件、查看分析结果、切换目录等菜单,并对非法输入做了鲁棒性处理。 + +### 控制台界面截图 + +完整功能包括:输入目录 → 展示文件列表 → 分析指定文件 → 分析全部文件 → 查看单个文件的解析结果。 + +![完整功能截图](./assets/localcli-full.png) + +### 鲁棒性测试截图 + +覆盖以下非法输入场景:不存在的目录、不存在的文件名、分析不存在的文件(抛出 `ArgumentException` 被捕获)、菜单选项输入非数字、以及查看尚未分析过的文件。 + +![鲁棒性测试截图](./assets/localcli-robustness.png) + +--- + +## Q2.1 + +### `WorkQueue` 中的共享变量及其保护 + +`WorkQueue` 中有两个共享变量: + ++ `_items`(`Queue`):队列内部存储; ++ `_isCompleted`(`bool`):标记是否已结束放入元素。 + +这两个变量均通过 `lock (_items)` 保护,即以 `_items` 对象本身作为互斥量。`Enqueue`、`TryDequeue`、`CompleteAdding` 以及 `IsCompleted` 属性中所有对这两个共享变量的读写都在 `lock (_items)` 临界区内完成,从而避免数据竞争。 + +### `LogFileAnalyzer` 中的共享变量及其保护 + +`LogFileAnalyzer` 中的共享变量有: + ++ `_currentDirectory`(`string?`):当前日志目录; ++ `_isAnalyzing`(`bool`):是否正在分析; ++ `_logFiles`(`Dictionary`):目录中的日志文件映射; ++ `_analysisResults`(`Dictionary`):各文件的解析结果。 + +这些变量统一通过一个专用的互斥对象 `_syncRoot` 的 `lock (_syncRoot)` 保护。所有方法(`ChangeDirectory`、`GetLogFiles`、`TryGetAnalysisResult`、`AnalyzeFiles`、`RunWorkers`、`WorkerMain` 等)在访问这些共享变量时都先进入 `lock (_syncRoot)` 临界区。尤其是 `_isAnalyzing` 的读写、`_analysisResults` 的读写(由多个工作线程同时写入),都必须加锁。 + +### 条件变量使用 `if` 而非 `while` 的后果 + +若将判断条件写成 `if`,当出现虚假唤醒(spurious wakeup)时:线程在没有人调用 `signal`/`broadcast` 的情况下从 `Monitor.Wait` 中醒来,但此时仓库(队列)可能仍然是空的。若用 `if`,线程醒来后不会再检查条件,而是直接越过等待、去执行取元素操作,这会导致: + ++ 从空队列中执行 `Dequeue()`,抛出 `InvalidOperationException`(或读到非法数据); ++ 对应到无限容量生产者-消费者问题,就是消费者在 `buffer == 0` 时依然执行 `buffer -= 1`,造成"负库存"的逻辑错误。 + +因此必须用 `while`,让线程每次被唤醒后都重新检查"是否有元素"以及"是否已结束放入",只有在条件真正满足时才继续执行,从而保证正确性。 + +--- + +## Q2.2 + +扫描目录中全部 `.log` 后缀日志文件的代码位于 `LogFileAnalyzer.ChangeDirectory` 方法中: + +```csharp +var logFiles = Directory.EnumerateFiles(directoryPath, "*.log", SearchOption.TopDirectoryOnly) + .Select(filePath => Path.GetFileName(filePath)) + .OrderBy(fileName => fileName); +``` + +它使用 `Directory.EnumerateFiles` 配合通配符 `"*.log"` 和 `SearchOption.TopDirectoryOnly` 枚举当前目录下的日志文件,再用 `Select` 取文件名、`OrderBy` 排序。 + +若要递归获取给定目录的全部子目录(及子子目录……)内的日志文件,只需把搜索选项 `SearchOption.TopDirectoryOnly` 改为 `SearchOption.AllDirectories` 即可(可配合 `Directory.EnumerateFiles(directoryPath, "*.log", SearchOption.AllDirectories)`)。 + +--- + +## Q2.3 + +### Q2.3.b + +本次作业我使用了 AI 辅助完成。我给予 AI 的提示词大致是:"切换到 02-multithreading 分支,阅读 guidance 文档后完成 WorkQueue、LogFileAnalyzer、LocalCli 的实现"。 + +我对 AI 的使用主要是:让 AI 帮我梳理 `WorkQueue` 的条件变量写法(`Monitor.Wait`/`Pulse`/`PulseAll` 与 `while` 循环配合)、`LogFileAnalyzer` 中 `RunWorkers` 的线程生命周期管理(入队 → 开启线程 → `Join`),以及 `LocalCli` 的异常捕获结构。 + +AI 的解答基本正确,但存在一些需要人工修正的细节:例如 `WorkQueue.TryDequeue` 中 `item = _items.Dequeue()` 会触发可空性警告 CS8762,需要用空值宽容运算符 `!` 修正;又如 `RunWorkers` 中"结束放入"(`CompleteAdding`)必须在开启工作线程之前(或之后立刻)执行,并唤醒所有消费者,否则消费者会永久阻塞在 `TryDequeue` 上。这些都是需要人工理解并发语义后自行确认的点。 + +我认为本节的难度为适中。 diff --git a/docs/03-async-grpc/assets/remote-cli-full.jpg b/docs/03-async-grpc/assets/remote-cli-full.jpg new file mode 100644 index 0000000..2c560f2 Binary files /dev/null and b/docs/03-async-grpc/assets/remote-cli-full.jpg differ diff --git a/docs/03-async-grpc/assets/remote-cli-robustness.jpg b/docs/03-async-grpc/assets/remote-cli-robustness.jpg new file mode 100644 index 0000000..2a8d6e5 Binary files /dev/null and b/docs/03-async-grpc/assets/remote-cli-robustness.jpg differ diff --git a/docs/03-async-grpc/report.md b/docs/03-async-grpc/report.md new file mode 100644 index 0000000..5f33c03 --- /dev/null +++ b/docs/03-async-grpc/report.md @@ -0,0 +1,50 @@ +# Report for Async and gRPC + +## 功能实现简介 + +本节在 `02-multithreading` 的基础上,实现了一个常驻运行的 gRPC Agent 服务,以及一个远程控制台客户端,共完成四个部分: + +1. **类型转换 `GrpcLogEntryVisitor` / `GrpcTypeConverter`**(`LogAnalyzerRpc`):利用访问者模式将 `LogEntry` 与 Protobuf 的 `LogEntryMessage` 互相转换,并完成 `LogSeverity`、`LogEventType`、`AnalysisState` 等枚举的相互转换。 +2. **gRPC 服务 `AgentService`**(`LogAnalyzerAgent/Services`):gRPC 服务入口,把用户请求转发给 `AgentSession`。 +3. **业务逻辑 `AgentSession`**(`LogAnalyzerAgent/Applications`):调用 `LogFileAnalyzer` 完成 `ChangeDirectory`、`GetLogFiles`、`AnalyzeAll`、`AnalyzeFiles`、`GetAnalysisResult`,并做好异常处理,保证服务永不因非法请求而崩溃。 +4. **远程控制台客户端 `RemoteCli`**:将上一节的 `LocalCli` 改装为 gRPC 客户端版本,全部使用异步调用,并具备非法输入鲁棒性。 + +### 完整功能截图 + +完整功能包括:启动 Agent → 运行 RemoteCli → 切换目录 → 展示文件列表 → 分析全部文件 → 查看单个文件的流式解析结果。 + +![完整功能截图](./assets/remote-cli-full.jpg) + +### 鲁棒性测试截图 + +覆盖场景:输入不存在的目录(返回 `DirectoryNotFound`)、查看不存在的文件(返回 `FileNotFound`)、查看尚未分析的文件(返回 `NotAnalyzed` 提示)、输入非法的并行度、菜单选项输入非数字等。 + +![鲁棒性测试截图](./assets/remote-cli-robustness.jpg) + +--- + +## Q3.1 + +我认为开发网络应用程序与以往开发非网络应用程序最大的区别在于**程序运行的边界不再局限于单个进程**,而在于: + +1. **错误来源更多、更不可控**:本地程序的错误基本可以确定(逻辑错误、越界等);而网络程序还需要面对网络抖动、对端进程崩溃、连接超时、序列化/反序列化失败、协议不匹配等大量额外的失败场景,必须假定"对端随时可能出错"。 +2. **需要处理跨语言、跨机器的一致性问题**:本地程序内存中的对象直接使用即可;网络程序则必须把内存中的对象序列化成协议(这里用 Protobuf),并在两端保持类型、字段、枚举的定义一致,任何一边改动都可能破坏通信。 +3. **并发与异步是常态**:网络程序天然是 I/O 密集型,阻塞式等待会浪费 CPU,因此需要异步编程(`async`/`await`)让线程在等待网络时不空转;同时服务端要同时响应多个客户端,必须考虑线程安全。 +4. **服务可用性是硬指标**:本地程序崩了重开即可;而 Agent 作为常驻服务,一旦因某个非法请求崩溃,所有客户端都无法使用,所以必须对所有输入做防御式处理,保证"尽可能不崩溃"。 +5. **调试更复杂**:需要同时运行服务端和客户端两个程序,还要观察网络上的实际请求/响应,比单进程调试繁琐得多。 + +额外的难点主要集中在:异常处理要覆盖网络层与业务层、数据一致性、并发安全,以及异步调用链的正确性。 + +--- + +## Q3.2 + +### Q3.2.b + +本次作业我使用了 AI 辅助完成。我给予 AI 的提示词大致是:"切换到 03-async-grpc 分支,阅读 guidance 文档后完成 GrpcLogEntryVisitor、GrpcTypeConverter、AgentSession、AgentService、RemoteCli 的实现"。 + +我对 AI 的使用主要是:询问 gRPC 客户端对 server-streaming 调用的接收方式(`client.GetAnalysisResult(request)` 返回 `AsyncServerStreamingCall`,用 `await call.ResponseStream.ReadAllAsync().ToListAsync()` 读取),以及 Protobuf `oneof` 与 `optional string` 在 C# 生成代码里的表现形式(`EntryOneofCase`/`PayloadOneofCase`、`HasErrorMessage` 等)。 + +AI 的解答基本正确,但有一些需要人工确认的细节:例如 `AnalysisResultHeaderMessage.ErrorMessage` 是 `optional string`,其 setter 会对 `null` 抛异常,因此需要先判断 `result.ErrorMessage is not null` 再赋值,否则分析成功的文件(`ErrorMessage == null`)会导致服务端出错;又如 Agent 必须捕获所有异常并转化为 `OperationStatusMessage`,不能直接把异常抛给客户端。这些都是需要结合生成代码与"服务不崩溃"这一目标自行判断的点。 + +从 AI 那里我还了解到:gRPC 是 lazy 连接(第一次调用时才真正建立连接,所以用 `Ping` 预热),以及服务端流式返回(`IServerStreamWriter`)与客户端流式读取的配对方式。整体而言,本节难度适中。 diff --git a/src/LocalCli/Program.cs b/src/LocalCli/Program.cs index 17b30db..f291fed 100644 --- a/src/LocalCli/Program.cs +++ b/src/LocalCli/Program.cs @@ -112,22 +112,92 @@ 6. Exit. private static void ShowLogFiles(LogFileAnalyzer analyzer) { - throw new NotImplementedException("T2.3"); + var logFiles = analyzer.GetLogFiles(); + if (logFiles.Count == 0) + { + Console.WriteLine("No log files found in the directory."); + return; + } + + Console.WriteLine("Log files:"); + foreach (var fileName in logFiles) + { + Console.WriteLine($" {fileName}"); + } } private static void AnalyzeFiles(LogFileAnalyzer analyzer) { - throw new NotImplementedException("T2.3"); + Console.WriteLine("Please input file names separated by commas (e.g., basic.log,basic-multiple.log):"); + var input = Console.ReadLine(); + if (input is null) + { + return; + } + + var fileNames = input.Split(',', StringSplitOptions.RemoveEmptyEntries | StringSplitOptions.TrimEntries); + if (fileNames.Length == 0) + { + Console.WriteLine("No file names provided, please try again."); + return; + } + + try + { + analyzer.AnalyzeFiles(0, fileNames); + Console.WriteLine("Analysis completed."); + } + catch (Exception ex) + { + Console.WriteLine($"Failed to analyze files: {ex.Message}"); + } } private static void AnalyzeAll(LogFileAnalyzer analyzer) { - throw new NotImplementedException("T2.3"); + try + { + analyzer.AnalyzeAll(0); + Console.WriteLine("Analysis completed."); + } + catch (Exception ex) + { + Console.WriteLine($"Failed to analyze all files: {ex.Message}"); + } } private static void GetAnalysisResult(LogFileAnalyzer analyzer) { - throw new NotImplementedException("T2.3"); + Console.WriteLine("Please input file name:"); + var fileName = Console.ReadLine(); + if (fileName is null) + { + return; + } + + if (!analyzer.TryGetAnalysisResult(fileName, out var result)) + { + Console.WriteLine($"File '{fileName}' is not in the current directory, please try again."); + return; + } + + switch (result!.State) + { + case AnalysisState.NotAnalyzed: + Console.WriteLine($"File '{fileName}' has not been analyzed yet."); + break; + case AnalysisState.Succeeded: + var visitor = new KeyValueVisitor(); + foreach (var entry in result.Entries) + { + var kv = visitor.Dump(entry); + Console.WriteLine(string.Join(", ", kv.Select(pair => $"{pair.Key}={pair.Value}"))); + } + break; + case AnalysisState.Failed: + Console.WriteLine($"File '{fileName}' failed to analyze: {result.ErrorMessage}"); + break; + } } } } diff --git a/src/LogAnalyzer/LogFileAnalyzer.cs b/src/LogAnalyzer/LogFileAnalyzer.cs index c3e7691..80d0f41 100644 --- a/src/LogAnalyzer/LogFileAnalyzer.cs +++ b/src/LogAnalyzer/LogFileAnalyzer.cs @@ -138,10 +138,7 @@ public void AnalyzeFiles(int degreeOfParallelism, IEnumerable fileNames) } fileList = fileNameList.Select(fileName => _logFiles[fileName]).ToList(); - /* - * Set _isAnalyzing - */ - // TODO: T2.2 + _isAnalyzing = true; } try @@ -150,11 +147,10 @@ public void AnalyzeFiles(int degreeOfParallelism, IEnumerable fileNames) } finally { - /* - * Unset _isAnalyzing - * Remember to lock _syncRoot to prevent data race - */ - // TODO: T2.2 + lock (_syncRoot) + { + _isAnalyzing = false; + } } } @@ -165,11 +161,14 @@ private void RunWorkers(int degreeOfParallelism, IReadOnlyList fileLis { foreach (var file in fileList) { - /* - * Filter unparsed files. - * If there is an unknown file, throw System.InvalidOperationException. - */ - throw new NotImplementedException("TODO: T2.2"); + if (!_analysisResults.TryGetValue(file.Name, out var analysisResult)) + { + throw new InvalidOperationException($"Unknown file: {file.Name}."); + } + if (analysisResult.State == AnalysisState.NotAnalyzed) + { + logFilesToParse.Add(file); + } } } @@ -180,10 +179,11 @@ private void RunWorkers(int degreeOfParallelism, IReadOnlyList fileLis var queue = new WorkQueue(); - /* - * Enqueue log files - */ - // TODO: T2.2 + foreach (var file in logFilesToParse) + { + queue.Enqueue(file); + } + queue.CompleteAdding(); degreeOfParallelism = Math.Max(Math.Min(degreeOfParallelism, logFilesToParse.Count), 1); var workers = new Thread[degreeOfParallelism]; @@ -191,16 +191,19 @@ private void RunWorkers(int degreeOfParallelism, IReadOnlyList fileLis { int workerId = i; string threadName = $"log-analyzer-worker-{workerId}"; - /* - * Create and start threads to run `WorkerMain` - */ - // TODO: T2.2 + var worker = new Thread(() => WorkerMain(workerId, queue)) + { + IsBackground = true, + Name = threadName, + }; + workers[i] = worker; + worker.Start(); } - /* - * Wait for (join) all threads to end - */ - // TODO: T2.2 + foreach (var worker in workers) + { + worker.Join(); + } } private void WorkerMain(int workerId, WorkQueue queue) @@ -212,20 +215,37 @@ private void WorkerMain(int workerId, WorkQueue queue) AnalysisResult result; try { - // Parse file - throw new NotImplementedException("TODO: T2.2"); + List entries; + using (var reader = new StreamReader(file.FullName)) + { + entries = parser.Parse(reader).ToList(); + } + + result = new AnalysisResult( + FileName: file.Name, + FullName: file.FullName, + State: AnalysisState.Succeeded, + Entries: entries, + ErrorMessage: null, + WorkerId: workerId + ); } catch (Exception ex) { - // Save exception message to result - throw new NotImplementedException("TODO: T2.2"); + result = new AnalysisResult( + FileName: file.Name, + FullName: file.FullName, + State: AnalysisState.Failed, + Entries: Array.Empty(), + ErrorMessage: ex.Message, + WorkerId: workerId + ); } - /* - * Save parse result. - * [!Important] Remember to lock _syncRoot to prevent data race. - */ - throw new NotImplementedException("TODO: T2.2"); + lock (_syncRoot) + { + _analysisResults[file.Name] = result; + } } } } diff --git a/src/LogAnalyzer/WorkQueue.cs b/src/LogAnalyzer/WorkQueue.cs index 23055a5..948c434 100644 --- a/src/LogAnalyzer/WorkQueue.cs +++ b/src/LogAnalyzer/WorkQueue.cs @@ -20,17 +20,40 @@ public bool IsCompleted public void Enqueue(T item) { - throw new NotImplementedException("TODO: T2.1"); + lock (_items) + { + _items.Enqueue(item); + Monitor.Pulse(_items); + } } public bool TryDequeue([NotNullWhen(true)] out T? item) { - throw new NotImplementedException("TODO: T2.1"); + lock (_items) + { + while (_items.Count == 0 && !_isCompleted) + { + Monitor.Wait(_items); + } + + if (_items.Count > 0) + { + item = _items.Dequeue()!; + return true; + } + + item = default; + return false; + } } public void CompleteAdding() { - throw new NotImplementedException("TODO: T2.1"); + lock (_items) + { + _isCompleted = true; + Monitor.PulseAll(_items); + } } } } diff --git a/src/LogAnalyzerAgent/Applications/AgentSession.cs b/src/LogAnalyzerAgent/Applications/AgentSession.cs index 2531f22..4043d54 100644 --- a/src/LogAnalyzerAgent/Applications/AgentSession.cs +++ b/src/LogAnalyzerAgent/Applications/AgentSession.cs @@ -38,6 +38,27 @@ private static OperationStatusMessage CreateNoErrorOperationStatus() }; } + private static OperationStatusMessage CreateErrorOperationStatus(AgentErrorCode code, string message) + { + return new OperationStatusMessage() + { + Success = false, + Code = code, + Message = message, + }; + } + + private static OperationStatusMessage CreateOperationStatusFromException(Exception ex) + { + return ex switch + { + ArgumentOutOfRangeException => CreateErrorOperationStatus(AgentErrorCode.InvalidArgument, ex.Message), + InvalidOperationException => CreateErrorOperationStatus(AgentErrorCode.InvalidOperation, ex.Message), + ArgumentException => CreateErrorOperationStatus(AgentErrorCode.FileNotFound, ex.Message), + _ => CreateErrorOperationStatus(AgentErrorCode.InternalError, ex.Message), + }; + } + public Task Ping(Empty empty, CancellationToken cancellationToken) { return Task.FromResult(new Empty()); @@ -79,22 +100,124 @@ public Task GetLogFiles(Empty empty, CancellationToken canc public Task ChangeDirectory(ChangeDirectoryRequest request, CancellationToken cancellationToken) { - throw new NotImplementedException("TODO: T3.1"); + var response = new ChangeDirectoryResponse(); + try + { + if (!_analyzer.ChangeDirectory(request.DirectoryPath)) + { + if (_analyzer.IsAnalyzing) + { + response.Status = CreateErrorOperationStatus( + AgentErrorCode.InvalidOperation, "Cannot change directory while the agent is analyzing."); + } + else + { + response.Status = CreateErrorOperationStatus( + AgentErrorCode.DirectoryNotFound, $"Directory not found: {request.DirectoryPath}"); + } + return Task.FromResult(response); + } + + response.Status = CreateNoErrorOperationStatus(); + response.CurrentDirectory = _analyzer.CurrentDirectory ?? ""; + response.FileNames.AddRange(_analyzer.GetLogFiles()); + } + catch (Exception ex) + { + response.Status = CreateInternalErrorOperationStatus(ex); + _logger.LogError(ex, "An error occurred while changing directory."); + } + return Task.FromResult(response); } public Task AnalyzeAll(AnalyzeAllRequest request, CancellationToken cancellationToken) { - throw new NotImplementedException("TODO: T3.1"); + var response = new AnalyzeAllResponse(); + try + { + _analyzer.AnalyzeAll(request.DegreeOfParallelism); + response.Status = CreateNoErrorOperationStatus(); + } + catch (Exception ex) + { + response.Status = CreateOperationStatusFromException(ex); + _logger.LogError(ex, "An error occurred while analyzing all log files."); + } + return Task.FromResult(response); } public Task AnalyzeFiles(AnalyzeFilesRequest request, CancellationToken cancellationToken) { - throw new NotImplementedException("TODO: T3.1"); + var response = new AnalyzeFilesResponse(); + try + { + _analyzer.AnalyzeFiles(request.DegreeOfParallelism, request.FileNames); + response.Status = CreateNoErrorOperationStatus(); + } + catch (Exception ex) + { + response.Status = CreateOperationStatusFromException(ex); + _logger.LogError(ex, "An error occurred while analyzing specified log files."); + } + return Task.FromResult(response); } public IReadOnlyList GetAnalysisResult(GetAnalysisResultRequest request, CancellationToken cancellationToken) { - throw new NotImplementedException("TODO: T3.1"); + var responses = new List(); + try + { + if (!_analyzer.TryGetAnalysisResult(request.FileName, out var result)) + { + responses.Add(new GetAnalysisResultResponse() + { + Status = CreateErrorOperationStatus( + AgentErrorCode.FileNotFound, $"File '{request.FileName}' is not in the current directory."), + }); + return responses; + } + + var header = new AnalysisResultHeaderMessage() + { + FileName = result!.FileName, + FullName = result.FullName, + State = GrpcTypeConverter.ConvertToGrpc(result.State), + WorkerId = result.WorkerId, + }; + if (result.ErrorMessage is not null) + { + header.ErrorMessage = result.ErrorMessage; + } + + responses.Add(new GetAnalysisResultResponse() + { + Header = header, + Status = CreateNoErrorOperationStatus(), + }); + + if (result.State == AnalysisState.Succeeded) + { + foreach (var entry in result.Entries) + { + responses.Add(new GetAnalysisResultResponse() + { + LogEntry = GrpcTypeConverter.ConvertToGrpc(entry), + Status = CreateNoErrorOperationStatus(), + }); + } + } + + return responses; + } + catch (Exception ex) + { + responses.Add(new GetAnalysisResultResponse() + { + Status = CreateInternalErrorOperationStatus(ex), + }); + _logger.LogError(ex, "An error occurred while retrieving analysis result."); + return responses; + } } } } diff --git a/src/LogAnalyzerAgent/Services/AgentService.cs b/src/LogAnalyzerAgent/Services/AgentService.cs index 591dcad..d38d1cf 100644 --- a/src/LogAnalyzerAgent/Services/AgentService.cs +++ b/src/LogAnalyzerAgent/Services/AgentService.cs @@ -29,27 +29,31 @@ public override Task GetAgentStatus(Empty empty, ServerC public override Task ChangeDirectory(ChangeDirectoryRequest request, ServerCallContext context) { - throw new NotImplementedException("TODO: T3.1"); + return _session.ChangeDirectory(request, context.CancellationToken); } public override Task GetLogFiles(Empty empty, ServerCallContext context) { - throw new NotImplementedException("TODO: T3.1"); + return _session.GetLogFiles(empty, context.CancellationToken); } public override Task AnalyzeAll(AnalyzeAllRequest request, ServerCallContext context) { - throw new NotImplementedException("TODO: T3.1"); + return _session.AnalyzeAll(request, context.CancellationToken); } public override Task AnalyzeFiles(AnalyzeFilesRequest request, ServerCallContext context) { - throw new NotImplementedException("TODO: T3.1"); + return _session.AnalyzeFiles(request, context.CancellationToken); } public override async Task GetAnalysisResult(GetAnalysisResultRequest request, IServerStreamWriter responseStream, ServerCallContext context) { - throw new NotImplementedException("TODO: T3.1"); + var responses = _session.GetAnalysisResult(request, context.CancellationToken); + foreach (var response in responses) + { + await responseStream.WriteAsync(response); + } } } } diff --git a/src/LogAnalyzerRpc/GrpcLogEntryVisitor.cs b/src/LogAnalyzerRpc/GrpcLogEntryVisitor.cs index eb69232..3196aac 100644 --- a/src/LogAnalyzerRpc/GrpcLogEntryVisitor.cs +++ b/src/LogAnalyzerRpc/GrpcLogEntryVisitor.cs @@ -30,12 +30,38 @@ public LogEntryMessage Visit(CallLogEntry entry) public LogEntryMessage Visit(RequestLogEntry entry) { - throw new NotImplementedException("TODO: T3.1"); + return new LogEntryMessage() + { + RequestLogEntry = new RequestLogEntryMessage + { + LineNo = entry.LineNo, + Timestamp = Timestamp.FromDateTimeOffset(entry.Timestamp), + PodName = entry.PodName, + Severity = GrpcTypeConverter.ConvertToGrpc(entry.Severity), + EventType = GrpcTypeConverter.ConvertToGrpc(entry.EventType), + RequestId = entry.RequestId, + Method = entry.Method, + Path = entry.Path, + StatusCode = entry.StatusCode, + } + }; } public LogEntryMessage Visit(InternalLogEntry entry) { - throw new NotImplementedException("TODO: T3.1"); + return new LogEntryMessage() + { + InternalLogEntry = new InternalLogEntryMessage + { + LineNo = entry.LineNo, + Timestamp = Timestamp.FromDateTimeOffset(entry.Timestamp), + PodName = entry.PodName, + Severity = GrpcTypeConverter.ConvertToGrpc(entry.Severity), + EventType = GrpcTypeConverter.ConvertToGrpc(entry.EventType), + ExceptionName = entry.ExceptionName, + ExceptionMessage = entry.ExceptionMessage, + } + }; } } } diff --git a/src/LogAnalyzerRpc/GrpcTypeConverter.cs b/src/LogAnalyzerRpc/GrpcTypeConverter.cs index 029122e..aa31134 100644 --- a/src/LogAnalyzerRpc/GrpcTypeConverter.cs +++ b/src/LogAnalyzerRpc/GrpcTypeConverter.cs @@ -20,12 +20,24 @@ public static AnalysisStateEnum ConvertToGrpc(AnalysisState state) public static LogSeverityEnum ConvertToGrpc(LogSeverity severity) { - throw new NotImplementedException("TODO: T3.1"); + return severity switch + { + LogSeverity.Info => LogSeverityEnum.Info, + LogSeverity.Warning => LogSeverityEnum.Warning, + LogSeverity.Error => LogSeverityEnum.Error, + _ => throw new ArgumentOutOfRangeException(nameof(severity), severity, null) + }; } public static LogEventTypeEnum ConvertToGrpc(LogEventType eventType) { - throw new NotImplementedException("TODO: T3.1"); + return eventType switch + { + LogEventType.Call => LogEventTypeEnum.Call, + LogEventType.Request => LogEventTypeEnum.Request, + LogEventType.Internal => LogEventTypeEnum.Internal, + _ => throw new ArgumentOutOfRangeException(nameof(eventType), eventType, null) + }; } public static LogEntryMessage ConvertToGrpc(LogEntry entry) @@ -46,12 +58,24 @@ public static AnalysisState ConvertFromGrpc(AnalysisStateEnum state) public static LogSeverity ConvertFromGrpc(LogSeverityEnum severity) { - throw new NotImplementedException("TODO: T3.1"); + return severity switch + { + LogSeverityEnum.Info => LogSeverity.Info, + LogSeverityEnum.Warning => LogSeverity.Warning, + LogSeverityEnum.Error => LogSeverity.Error, + _ => throw new ArgumentOutOfRangeException(nameof(severity), severity, null) + }; } public static LogEventType ConvertFromGrpc(LogEventTypeEnum eventType) { - throw new NotImplementedException("TODO: T3.1"); + return eventType switch + { + LogEventTypeEnum.Call => LogEventType.Call, + LogEventTypeEnum.Request => LogEventType.Request, + LogEventTypeEnum.Internal => LogEventType.Internal, + _ => throw new ArgumentOutOfRangeException(nameof(eventType), eventType, null) + }; } public static LogEntry ConvertFromGrpc(LogEntryMessage entryMessage) @@ -67,8 +91,24 @@ public static LogEntry ConvertFromGrpc(LogEntryMessage entryMessage) TargetService: entryMessage.CallLogEntry.TargetService, DurationMs: entryMessage.CallLogEntry.DurationMs ), - LogEntryMessage.EntryOneofCase.RequestLogEntry => throw new NotImplementedException("TODO: T3.1"), - LogEntryMessage.EntryOneofCase.InternalLogEntry => throw new NotImplementedException("TODO: T3.1"), + LogEntryMessage.EntryOneofCase.RequestLogEntry => new RequestLogEntry( + LineNo: entryMessage.RequestLogEntry.LineNo, + Timestamp: entryMessage.RequestLogEntry.Timestamp.ToDateTimeOffset(), + PodName: entryMessage.RequestLogEntry.PodName, + Severity: ConvertFromGrpc(entryMessage.RequestLogEntry.Severity), + RequestId: entryMessage.RequestLogEntry.RequestId, + Method: entryMessage.RequestLogEntry.Method, + Path: entryMessage.RequestLogEntry.Path, + StatusCode: entryMessage.RequestLogEntry.StatusCode + ), + LogEntryMessage.EntryOneofCase.InternalLogEntry => new InternalLogEntry( + LineNo: entryMessage.InternalLogEntry.LineNo, + Timestamp: entryMessage.InternalLogEntry.Timestamp.ToDateTimeOffset(), + PodName: entryMessage.InternalLogEntry.PodName, + Severity: ConvertFromGrpc(entryMessage.InternalLogEntry.Severity), + ExceptionName: entryMessage.InternalLogEntry.ExceptionName, + ExceptionMessage: entryMessage.InternalLogEntry.ExceptionMessage + ), _ => throw new ArgumentException($"Unknown entry type: {entryMessage.EntryCase}", nameof(entryMessage)) }; } diff --git a/src/LogParser/Models/LogEntries.cs b/src/LogParser/Models/LogEntries.cs index 69edbc0..e4e9bbc 100644 --- a/src/LogParser/Models/LogEntries.cs +++ b/src/LogParser/Models/LogEntries.cs @@ -54,7 +54,7 @@ public sealed record RequestLogEntry( { public override TResult Accept(ILogEntryVisitor visitor) { - throw new NotImplementedException("TODO: T1.2"); + return visitor.Visit(this); } } @@ -69,7 +69,7 @@ public sealed record InternalLogEntry( { public override TResult Accept(ILogEntryVisitor visitor) { - throw new NotImplementedException("TODO: T1.2"); + return visitor.Visit(this); } } diff --git a/src/LogParser/Parser/LineParser.cs b/src/LogParser/Parser/LineParser.cs index 0475f6b..6b19485 100644 --- a/src/LogParser/Parser/LineParser.cs +++ b/src/LogParser/Parser/LineParser.cs @@ -16,8 +16,8 @@ public static LogEntry ParseLine(LogRecord logRecord) return eventElement.GetString() switch { "call" => LineParser.CreateCall(logRecord), - "request" => throw new NotImplementedException("TODO: T1.2"), - "internal" => throw new NotImplementedException("TODO: T1.2"), + "request" => LineParser.CreateRequest(logRecord), + "internal" => LineParser.CreateInternal(logRecord), _ => throw new FormatException($"Unknown event type: {eventElement.GetString()} in log message: {logRecord.Message}") }; } @@ -50,12 +50,39 @@ private static LogEntry CreateCall(LogRecord logRecord) private static LogEntry CreateRequest(LogRecord logRecord) { - throw new NotImplementedException("TODO: T1.2"); + var requestMessage = JsonSerializer.Deserialize(logRecord.Message, options) + ?? throw new FormatException($"Failed to deserialize request message: {logRecord.Message}"); + return new RequestLogEntry( + LineNo: logRecord.LineNo, + Timestamp: DateTimeOffset.Parse(logRecord.Timestamp), + PodName: logRecord.PodName, + Severity: ParseSeverity(requestMessage.Severity), + RequestId: requestMessage.RequestId, + Method: requestMessage.Method, + Path: requestMessage.Path, + StatusCode: requestMessage.StatusCode + ); } private static LogEntry CreateInternal(LogRecord logRecord) { - throw new NotImplementedException("TODO: T1.2"); + var internalMessage = JsonSerializer.Deserialize(logRecord.Message, options) + ?? throw new FormatException($"Failed to deserialize internal message: {logRecord.Message}"); + var separatorIndex = internalMessage.Exception.IndexOf(": "); + if (separatorIndex < 0) + { + throw new FormatException($"Invalid exception format: {internalMessage.Exception}"); + } + var exceptionName = internalMessage.Exception[..separatorIndex]; + var exceptionMessage = internalMessage.Exception[(separatorIndex + 2)..]; + return new InternalLogEntry( + LineNo: logRecord.LineNo, + Timestamp: DateTimeOffset.Parse(logRecord.Timestamp), + PodName: logRecord.PodName, + Severity: ParseSeverity(internalMessage.Severity), + ExceptionName: exceptionName, + ExceptionMessage: exceptionMessage + ); } private static LogSeverity ParseSeverity(string severity) @@ -77,11 +104,16 @@ private record CallMessage( ); private record RequestMessage( - // TODO: T1.2 + [property: JsonRequired] string Severity, + [property: JsonRequired] string RequestId, + [property: JsonRequired] string Method, + [property: JsonRequired] string Path, + [property: JsonRequired] int StatusCode ); private record InternalMessage( - // TODO: T1.2 + [property: JsonRequired] string Severity, + [property: JsonRequired] string Exception ); } } diff --git a/src/LogParser/Visitors/KeyValueVisitor.cs b/src/LogParser/Visitors/KeyValueVisitor.cs index e5ceba2..f70bcc2 100644 --- a/src/LogParser/Visitors/KeyValueVisitor.cs +++ b/src/LogParser/Visitors/KeyValueVisitor.cs @@ -26,12 +26,32 @@ public Dictionary Visit(CallLogEntry entry) public Dictionary Visit(RequestLogEntry entry) { - throw new NotImplementedException("TODO: T1.3"); + return new Dictionary + { + ["LineNo"] = entry.LineNo.ToString(), + ["Timestamp"] = entry.Timestamp.ToString("O"), + ["PodName"] = entry.PodName, + ["Severity"] = entry.Severity.ToString(), + ["EventType"] = entry.EventType.ToString(), + ["RequestId"] = entry.RequestId, + ["Method"] = entry.Method, + ["Path"] = entry.Path, + ["StatusCode"] = entry.StatusCode.ToString(), + }; } public Dictionary Visit(InternalLogEntry entry) { - throw new NotImplementedException("TODO: T1.3"); + return new Dictionary + { + ["LineNo"] = entry.LineNo.ToString(), + ["Timestamp"] = entry.Timestamp.ToString("O"), + ["PodName"] = entry.PodName, + ["Severity"] = entry.Severity.ToString(), + ["EventType"] = entry.EventType.ToString(), + ["ExceptionName"] = entry.ExceptionName, + ["ExceptionMessage"] = entry.ExceptionMessage, + }; } } } diff --git a/src/RemoteCli/Program.cs b/src/RemoteCli/Program.cs index de0ac99..5999840 100644 --- a/src/RemoteCli/Program.cs +++ b/src/RemoteCli/Program.cs @@ -116,32 +116,160 @@ 6. Exit. private static async Task ShowLogFiles(LogAnalyzerAgentServiceClient client) { - throw new NotImplementedException("TODO: T3.2"); + var response = await client.GetLogFilesAsync(new Empty()); + if (!response.Status.Success) + { + Console.WriteLine($"Error: {response.Status.Code}: {response.Status.Message}"); + return; + } + + if (response.FileNames.Count == 0) + { + Console.WriteLine("No log files found in the directory."); + return; + } + + Console.WriteLine("Log files:"); + foreach (var fileName in response.FileNames) + { + Console.WriteLine($" {fileName}"); + } } private static int ReadDegreeOfParallelism() { - throw new NotImplementedException("TODO: T3.2"); + while (true) + { + Console.WriteLine("Please input degree of parallelism (0 for auto):"); + var input = Console.ReadLine(); + if (input is null) + { + return 0; + } + if (int.TryParse(input, out var degree) && degree >= 0) + { + return degree; + } + Console.WriteLine("Invalid degree of parallelism, please try again."); + } } private static List ReadFileNames() { - throw new NotImplementedException("TODO: T3.2"); + Console.WriteLine("Please input file names separated by commas (e.g., basic.log,basic-multiple.log):"); + var input = Console.ReadLine(); + if (input is null) + { + return new List(); + } + + var fileNames = input.Split(',', StringSplitOptions.RemoveEmptyEntries | StringSplitOptions.TrimEntries).ToList(); + if (fileNames.Count == 0) + { + Console.WriteLine("No file names provided."); + } + return fileNames; } private static async Task AnalyzeFiles(LogAnalyzerAgentServiceClient client) { - throw new NotImplementedException("TODO: T3.2"); + var degree = ReadDegreeOfParallelism(); + var fileNames = ReadFileNames(); + if (fileNames.Count == 0) + { + return; + } + + var request = new AnalyzeFilesRequest() + { + DegreeOfParallelism = degree, + }; + request.FileNames.AddRange(fileNames); + + var response = await client.AnalyzeFilesAsync(request); + if (response.Status.Success) + { + Console.WriteLine("Analysis completed."); + } + else + { + Console.WriteLine($"Error: {response.Status.Code}: {response.Status.Message}"); + } } private static async Task AnalyzeAll(LogAnalyzerAgentServiceClient client) { - throw new NotImplementedException("TODO: T3.2"); + var degree = ReadDegreeOfParallelism(); + var response = await client.AnalyzeAllAsync(new AnalyzeAllRequest() + { + DegreeOfParallelism = degree, + }); + + if (response.Status.Success) + { + Console.WriteLine("Analysis completed."); + } + else + { + Console.WriteLine($"Error: {response.Status.Code}: {response.Status.Message}"); + } } private static async Task GetAnalysisResult(LogAnalyzerAgentServiceClient client) { - throw new NotImplementedException("TODO: T3.2"); + Console.WriteLine("Please input file name:"); + var fileName = Console.ReadLine(); + if (fileName is null) + { + return; + } + + var request = new GetAnalysisResultRequest() + { + FileName = fileName, + }; + + using var call = client.GetAnalysisResult(request); + var responses = await call.ResponseStream.ReadAllAsync().ToListAsync(); + + if (responses.Count == 0) + { + Console.WriteLine("No response received."); + return; + } + + var first = responses[0]; + if (!first.Status.Success) + { + Console.WriteLine($"Error: {first.Status.Code}: {first.Status.Message}"); + return; + } + + if (first.PayloadCase != GetAnalysisResultResponse.PayloadOneofCase.Header) + { + Console.WriteLine("Unexpected response."); + return; + } + + var header = first.Header; + switch (header.State) + { + case AnalysisStateEnum.NotAnalyzed: + Console.WriteLine($"File '{fileName}' has not been analyzed yet."); + break; + case AnalysisStateEnum.Failed: + Console.WriteLine($"File '{fileName}' failed to analyze: {header.ErrorMessage}"); + break; + case AnalysisStateEnum.Succeeded: + var visitor = new KeyValueVisitor(); + foreach (var response in responses.Skip(1)) + { + var entry = GrpcTypeConverter.ConvertFromGrpc(response.LogEntry); + var kv = visitor.Dump(entry); + Console.WriteLine(string.Join(", ", kv.Select(pair => $"{pair.Key}={pair.Value}"))); + } + break; + } } } }