diff --git a/docs/01-basic/report.md b/docs/01-basic/report.md new file mode 100644 index 0000000..3e85fe2 --- /dev/null +++ b/docs/01-basic/report.md @@ -0,0 +1,60 @@ +# 报告 + +## Q1.1 的回答 + ++ 哪条语句或哪几条语句将日志按逗号进行分割?代码中,我们是如何指定每一行的第几个字段代表何种意义的? + - "分割"的实现位于 `LogFileParser.cs` 的 38 行 `foreach (var logRecord in csv.GetRecords())` 中实现,该行调用了 csv 库的 `GetRecords` 方法,将日志转化为我们所需要的 `LogRecord` 类。 + - 我们通过该文件 20-23 行 + + ```csharp + Map(m => m.LineNo).Index(0); + Map(m => m.Timestamp).Index(1); + Map(m => m.PodName).Index(2); + Map(m => m.Message).Index(3); + ``` + + 确定了:第一个字段是 `LineNo`,第二个字段是 `Timestamp`,以此类推。 + ++ 在对日志中 JSON 格式的 `message` 字段进行读取时,我们是在哪个方法内用哪几条语句判断这一行日志的种类(Call / Request / Internal)的? + - 我们在 `LineParser.cs` 的 `ParseLine` 方法中判断种类,具体来说,我们试图通过 + + ```csharp + 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($"Unknown event type: {eventElement.GetString()} in log message: {logRecord.Message}") + }; + } + ``` + + 检测传入 `LogRecord` 对象的 `message` 属性的 `event` 字段是三者中的哪一个,从而进行判断。 + ++ 在确定了日志种类后,我们是调用了哪个库方法对 JSON 进行解析的? + + 进一步,我们的框架代码是如何防止日志中有字段缺失的?(例如所给的 Call 日志的 `message` 中缺失 `request_id` 字段) + + 更进一步,日志中的 JSON 的键是 `abc-def` 命名法(称为烤串命名法),而我们的解析结果却是放在 `AbcDef` 命名法(称为大驼峰命名法)的属性里,我们的框架代码中是如何告诉 JSON 解析器完成这一命名法转换的? + - 我们调用了 `JsonSerializer` 的 `Deserialize` 方法解析 JSON。 + - 我们首先通过 `[property: JsonRequired]` 来确保:如果缺失字段,则抛出 `JsonException` 异常。此外,我们还通过 `??` 运算符检测 `JsonSerializer.Deserialize` 方法返回值是否为 `null`,如果在某些情况下该方法返回了 `null`,则抛出 `FormatException` 异常。 + - 我们在 `LineParser.cs` 的 31-34 行将 `options` 设为具有 `PropertyNamingPolicy = JsonNamingPolicy.KebabCaseLower` 的 `JsonSerializerOptions`,并在 `Deserialize` 时传入 `options` 参数,从而完成了命名法转换。 + +## Q1.2 的回答 + ++ `Dictionary KeyValueVisitor.Dump(LogEntry entry)` ++ `TResult Accept(ILogEntryVisitor visitor);` ++ `Dictionary Visit(CallLogEntry entry)` + +## Q1.3 的回答 + ++ 本次作业中,你是否使用了 AI? + - 我使用了 AI。 + +### Q1.3.b 的回答 + ++ 如果使用了 AI,你给予 AI 的提示词是什么?你认为 AI 给出的解答、你完全凭借传统搜索引擎以及自己的能力能够写出的解答之间,AI 的解答比你好在哪?AI 又有哪些解答是存在问题的,或者至少是不如你自己的解答的?给出你的理由。 + - 我主要使用的是 VS Code Copilot 自带的代码补全 AI,故没有给出提示词。 + - AI 的解答相比于自己给出的解答更为安全,且可读性更好,例如:`LineParser.cs` 的 72-75 这几行就是 AI 补充的,防止出现没有冒号的情况。 + - AI 仅仅看到了 `InternalMessage` 上文的两个格式就开始了 `private record InternalMessage` 的编写,但是它没有考虑到原始文本的 `ExceptionName` 和 `ExceptionMessage` 并不是由 JSON 解析给出的,我在测试未通过后,通过检查解决了该问题。 + - 本报告完全由我所写,但是由于我对于markdown的格式并不熟悉,所以我让AI调了一下格式 diff --git a/docs/02-multithreading/assets/localcli_func.png b/docs/02-multithreading/assets/localcli_func.png new file mode 100644 index 0000000..98e0500 Binary files /dev/null and b/docs/02-multithreading/assets/localcli_func.png differ diff --git a/docs/02-multithreading/assets/localcli_robust.png b/docs/02-multithreading/assets/localcli_robust.png new file mode 100644 index 0000000..4ae2496 Binary files /dev/null and b/docs/02-multithreading/assets/localcli_robust.png differ diff --git a/docs/02-multithreading/report.md b/docs/02-multithreading/report.md new file mode 100644 index 0000000..972f7f0 --- /dev/null +++ b/docs/02-multithreading/report.md @@ -0,0 +1,127 @@ +# Report: LocalCli Console Interface (T2.3) + +## 实现功能 + +根据 [guidance.md](./guidance.md) 中 Task 2.3 (S2.3) 的要求,完成了 `LocalCli/Program.cs` 中的控制台交互界面,包含以下功能: + +### 1. `InputDirectory` — 输入日志目录 + +提示用户输入日志文件所在目录,调用 `LogFileAnalyzer` 构造器扫描 `.log` 文件: + +- 目录不存在时调用 `analyzer.ChangeDirectory()` 返回 `false`,提示 "Directory not exists" 并要求重试 +- 目录路径非法(如空字符串)时捕获 `ArgumentException`,提示 "Directory illegal" 并要求重试 +- 用户输入 `Ctrl+C` / `Ctrl+Z`(`Console.ReadLine()` 返回 `null`)时安全退出 + +### 2. `ShowLogFiles` — 显示日志文件列表 + +调用 `analyzer.GetLogFiles()` 获取目录中所有 `.log` 文件,逐行打印文件名。 + +### 3. `AnalyzeFiles` — 分析指定日志文件 + +- 调用 `ReadDegreeOfParallelism()` 读取并行度(0 = 自动 / 逻辑处理器数),非数字或负数会提示重新输入 +- 调用 `ReadFileNames()` 读取逗号分隔的文件名列表(自动 `Trim` 并去除空项) +- 调用 `analyzer.AnalyzeFiles(degreeOfParallelism, fileNames)` 进行分析 +- 异常被捕获并以 "分析失败" 提示,程序不会崩溃 + +### 4. `AnalyzeAll` — 分析全部日志文件 + +- 读取并行度后调用 `analyzer.AnalyzeAll(degreeOfParallelism)` +- 同样做了异常捕获以保证鲁棒性 + +### 5. `GetAnalysisResult` — 获取分析结果 + +输入文件名,调用 `analyzer.TryGetAnalysisResult()`,分四种情况处理: + +| 情况 | 行为 | +|------|------| +| 文件不存在 | 提示 "File 'xxx' not found." | +| 尚未分析 (`NotAnalyzed`) | 提示 "File 'xxx' has not been analyzed." | +| 分析成功 (`Succeeded`) | 调用 `KeyValueVisitor.Dump` 逐行输出键值对 | +| 分析失败 (`Failed`) | 输出 `result.ErrorMessage` | + +### 6. 鲁棒性设计 + +对所有异常输入均有处理,程序不会崩溃: + +- 非法目录 → 提示重试 +- 非法菜单选项(非数字、超出范围)→ 提示 "Invalid choice/input" +- 非法并行度(负数、非数字)→ 提示重试 +- 空文件名列表 → 提示 "No file names input." 并返回主菜单 +- 不存在 / 未分析的文件查结果 → 给出明确提示 +- 切换目录后重新分析 → 结果正常重置 + +--- + +## 功能测试截图 + +### 完整功能演示 + +![完整功能演示](./assets/localcli_func.png) + +以上截图展示了完整的功能流程: +1. 输入目录 `dataset` +2. 显示日志文件列表(选项 1) +3. 分析指定文件 `basic.log, basic-fail.log`,并行度 2(选项 2) +4. 查看 `basic.log` 解析成功的 3 条记录(选项 4) +5. 查看 `basic-fail.log` 解析失败的错误信息(选项 4) +6. 分析全部文件(选项 3,并行度 0 = auto) +7. 查看 `basic-multiple.log` 200 条解析结果(选项 4) + +### 鲁棒性测试 + +![鲁棒性测试](./assets/localcli_robust.png) + +以上截图展示了各种非法输入的处理: +1. 输入不存在的目录 → 提示 "Directory not exists" 并重试 +2. 非法菜单选项 `0`, `abc`, `7` → 提示 "Invalid" +3. 非法并行度 `abc`, `-1` → 提示 "Invalid input" +4. 空文件名 → 提示 "No file names input." +5. 查不存在的文件的结果 → "File 'nonexistent.log' not found." +6. 查未分析文件的结果 → "has not been analyzed." +7. 分析全部 → 查看 `basic-fail.log` 的失败信息 +8. 切换目录后结果重置 → `basic-fail.log` 变回 "not been analyzed" +9. 重新分析后再次查看失败文件 → 正确输出错误信息 + +--- + +## 问答 + +### Q2.1 + +**`WorkQueue` 类中的共享变量有哪些?是通过什么保护其免于数据竞争(data race)呢?** + +`_items` 和 `_isCompleted` 都是共享变量,访问二者时提前打上 `_items` 的锁使得二者免于数据竞争。 + +**`LogFileAnalyzer` 类中的共享变量有哪些?是通过什么保护其免于数据竞争呢?** + +```csharp +private string? _currentDirectory = null; +private bool _isAnalyzing = false; +private readonly Dictionary _logFiles = new(); +private readonly Dictionary _analysisResults = new(); +``` + +以上都是共享变量,通过打上 `_syncRoot` 的锁避免数据竞争。 + +**如果条件变量的判断条件使用了 `if` 判断而非 `while` 判断,当出现了虚假唤醒现象时(在类 UNIX 系统中,由于 UNIX 信号等机制,即使没有人调用过 `signal` 或 `broadcast`,处于 `wait` 当中的条件变量也可能被唤醒),会出现什么后果?结合无限仓库容量的生产者消费者问题简单叙述一下。** + +以无限仓库容量的生产者消费者问题为例,如果采用 `if(queue.Count == 0)`,当内部 wait 虚假唤醒,则线程继续执行下方 `Dequeue`,试图出队空队列,从而导致抛出 `InvalidOperationException`。 + +### Q2.2 + +**那一段代码扫描了给定的目录中的全部 `.log` 后缀的日志文件?假使给定的需求是不但要扫描给定目录中的日志文件,还要递归地获取给定的目录的全部子目录、子子目录……内的日志文件,应当如何做(简要回答即可)?** + +```csharp +var logFiles = Directory.EnumerateFiles(directoryPath, "*.log", SearchOption.TopDirectoryOnly) + .Select(filePath => Path.GetFileName(filePath)) + .OrderBy(fileName => fileName); +``` + +以上代码扫描全部 `.log` 后缀的日志文件。将上述 `TopDirectoryOnly` 改为 `AllDirectories` 即可递归获取所有子目录中的日志文件。 + +### Q2.3 + +- 我使用了AI工具辅助 +- 第一次:由于我不会写文件的流式读取,导致 `parser.Parse` 参数类型不匹配,因此我在 VS Code 的 CC 插件中询问如下问题:"这段代码中 `result = parser.Parse(file);` 并不正确,`parser.Parse` 需要 `TextReader` 类型,应当如何修改?" +- 第二次:我借助了AI完成CLI:提示词为"根据 `docs/02-multiheading/guidance.md` 中对于 T2.3 的要求,完成 `Program.cs`" +- 本文件的测试部分也由AI生成,经过核对与 `report.md` 中要求相符 diff --git a/docs/03-async-grpc/assets/image.png b/docs/03-async-grpc/assets/image.png new file mode 100644 index 0000000..46b4eb5 Binary files /dev/null and b/docs/03-async-grpc/assets/image.png differ diff --git a/docs/03-async-grpc/assets/image_1.png b/docs/03-async-grpc/assets/image_1.png new file mode 100644 index 0000000..5922a37 Binary files /dev/null and b/docs/03-async-grpc/assets/image_1.png differ diff --git a/docs/03-async-grpc/report.md b/docs/03-async-grpc/report.md new file mode 100644 index 0000000..4be8628 --- /dev/null +++ b/docs/03-async-grpc/report.md @@ -0,0 +1,33 @@ +# 异步 gRPC 实验报告 + +## T3.2 截图与功能实现 + +### 输入图片 + +![输入图片](./assets/image.png) + +### 输出图片 + +![输出图片](./assets/image_1.png) + +## Q3.1 + +我认为,在 gRPC 的帮助下,二者的区别被减小了很多。但是以往我“开发”的应用程序往往都是自己使用或者测试,不太需要考虑特别高度的稳定性与对于输入的鲁棒性,因为即使程序崩掉了,也不过是重启一下,但是在本次的作业里,我发现了大量的异常处理等等保证鲁棒性的措施。此外,以往的开发中,我面对的往往是 CPU 密集型任务,因此不太常用多线程,但是网络中,通信成为了一个障碍,所以我在本次作业中也学着使用了一些异步语法。 + +我认为,本次的网络开发的主要难点在于,之前我的思维中,“服务”与“客户”是耦合的,但是在本次作业中,我们需要对二者做高度的解耦,这给我的理解带来了困难。 + +我认为本次主要的复杂之处在于信息流的长度比以前长了很多,以前简单的类型现在可能需要从服务端编码为 Protobuf 再在客户端解码,我认为这是较为复杂的。 + +## Q3.2 + +本次我使用了 AI。 + +### Q3.2.b + +**提示词:** + +> 这段代码中的最后一个方法应当实现流式返回,但是没有传入流对象,我应当如何处理? + +这是因为我不理解为什么 `AgentSession.GetAnalysisResult` 没有传入流对象,但是教学里对于流式 gRPC 要求传入一个流对象,AI 为我解答:这个方法并不是 gRPC 服务方法本身,其本身在 `AgentService` 里,我又查看了 `AgentService` 的代码,之后我理解了其层级。 + +由于我不会 Markdown 语法,本文的 Markdown 格式由 AI 进行了修改。 diff --git a/src/LocalCli/Program.cs b/src/LocalCli/Program.cs index 17b30db..25412e2 100644 --- a/src/LocalCli/Program.cs +++ b/src/LocalCli/Program.cs @@ -112,22 +112,115 @@ 6. Exit. private static void ShowLogFiles(LogFileAnalyzer analyzer) { - throw new NotImplementedException("T2.3"); + var _logfiles = analyzer.GetLogFiles(); + foreach (var file in _logfiles) + { + Console.WriteLine(file); + } + } + + private static int ReadDegreeOfParallelism() + { + while (true) + { + Console.WriteLine("Please input the degree of parallelism (0 means auto):"); + Console.Write(">>> "); + Console.Out.Flush(); + var str = Console.ReadLine(); + if (str is null) + { + return 0; + } + if (int.TryParse(str, out var degree) && degree >= 0) + { + return degree; + } + Console.WriteLine("Invalid input, please try again."); + } + } + + private static List ReadFileNames() + { + Console.WriteLine("Please input file names to analyze, separated by commas:"); + Console.Write(">>> "); + Console.Out.Flush(); + var str = Console.ReadLine(); + if (str is null) + { + return []; + } + return [.. str.Split(',', StringSplitOptions.TrimEntries | StringSplitOptions.RemoveEmptyEntries)]; } private static void AnalyzeFiles(LogFileAnalyzer analyzer) { - throw new NotImplementedException("T2.3"); + var degreeOfParallelism = ReadDegreeOfParallelism(); + var fileNames = ReadFileNames(); + if (fileNames.Count == 0) + { + Console.WriteLine("No file names input."); + return; + } + + try + { + analyzer.AnalyzeFiles(degreeOfParallelism, fileNames); + Console.WriteLine("Analysis finished."); + } + catch (Exception ex) + { + Console.WriteLine($"Analysis failed: {ex.Message}"); + } } private static void AnalyzeAll(LogFileAnalyzer analyzer) { - throw new NotImplementedException("T2.3"); + var degreeOfParallelism = ReadDegreeOfParallelism(); + try + { + analyzer.AnalyzeAll(degreeOfParallelism); + Console.WriteLine("Analysis finished."); + } + catch (Exception ex) + { + Console.WriteLine($"Analysis failed: {ex.Message}"); + } } private static void GetAnalysisResult(LogFileAnalyzer analyzer) { - throw new NotImplementedException("T2.3"); + Console.WriteLine("Please input the file name:"); + Console.Write(">>> "); + Console.Out.Flush(); + var fileName = Console.ReadLine(); + if (fileName is null) + { + return; + } + + if (!analyzer.TryGetAnalysisResult(fileName, out var result)) + { + Console.WriteLine($"File '{fileName}' not found."); + return; + } + + switch (result!.State) + { + case AnalysisState.NotAnalyzed: + Console.WriteLine($"File '{fileName}' has not been analyzed."); + break; + case AnalysisState.Succeeded: + var dumper = new KeyValueVisitor(); + foreach (var entry in result.Entries) + { + var kvPairs = dumper.Dump(entry); + Console.WriteLine(string.Join(", ", kvPairs.Select(kv => $"{kv.Key}: {kv.Value}"))); + } + break; + case AnalysisState.Failed: + Console.WriteLine($"Analysis of file '{fileName}' failed: {result.ErrorMessage}"); + break; + } } } } diff --git a/src/LogAnalyzer/LogFileAnalyzer.cs b/src/LogAnalyzer/LogFileAnalyzer.cs index c3e7691..47ea9d1 100644 --- a/src/LogAnalyzer/LogFileAnalyzer.cs +++ b/src/LogAnalyzer/LogFileAnalyzer.cs @@ -138,10 +138,8 @@ public void AnalyzeFiles(int degreeOfParallelism, IEnumerable fileNames) } fileList = fileNameList.Select(fileName => _logFiles[fileName]).ToList(); - /* - * Set _isAnalyzing - */ - // TODO: T2.2 + _isAnalyzing = true; + Monitor.PulseAll(_syncRoot); } try @@ -150,11 +148,11 @@ 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; + Monitor.PulseAll(_syncRoot); + } } } @@ -165,12 +163,20 @@ 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"); + AnalysisResult? result; + if(_analysisResults.TryGetValue(file.Name, out result)) + { + if(result != null && result.State == AnalysisState.NotAnalyzed) + { + logFilesToParse.Add(file); + } + } + else + { + throw new InvalidOperationException($"File '{file.Name}' is not in the current directory or does not exist."); + } } + Monitor.PulseAll(_syncRoot); } if (logFilesToParse.Count == 0) @@ -180,27 +186,34 @@ 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]; for (int i = 0; i < degreeOfParallelism; i++) { int workerId = i; string threadName = $"log-analyzer-worker-{workerId}"; - /* - * Create and start threads to run `WorkerMain` - */ - // TODO: T2.2 + workers[i] = new Thread(() => WorkerMain(workerId, queue)) + { + Name = threadName, + IsBackground = true + }; } - /* - * Wait for (join) all threads to end - */ - // TODO: T2.2 + foreach (var worker in workers) + { + worker.Start(); + } + + foreach (var worker in workers) + { + worker.Join(); + } } private void WorkerMain(int workerId, WorkQueue queue) @@ -212,20 +225,34 @@ private void WorkerMain(int workerId, WorkQueue queue) AnalysisResult result; try { - // Parse file - throw new NotImplementedException("TODO: T2.2"); + using var reader = file.OpenText(); + var 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; + Monitor.PulseAll(_syncRoot); + } } } } diff --git a/src/LogAnalyzer/WorkQueue.cs b/src/LogAnalyzer/WorkQueue.cs index 23055a5..c565ed5 100644 --- a/src/LogAnalyzer/WorkQueue.cs +++ b/src/LogAnalyzer/WorkQueue.cs @@ -20,17 +20,41 @@ public bool IsCompleted public void Enqueue(T item) { - throw new NotImplementedException("TODO: T2.1"); + lock(_items){ + if (_isCompleted){ + throw new InvalidOperationException("Cannot enqueue a completed queue."); + } + _items.Enqueue(item); + Monitor.PulseAll(_items); + } } public bool TryDequeue([NotNullWhen(true)] out T? item) { - throw new NotImplementedException("TODO: T2.1"); + lock (_items) + { + while(_items.Count == 0) + { + if (_isCompleted) + { + item = default; + return false; + } + Monitor.Wait(_items); + } + item = _items.Dequeue() ?? throw new InvalidOperationException("Queue is empty."); + Monitor.PulseAll(_items); + return true; + } } 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..8736423 100644 --- a/src/LogAnalyzerAgent/Applications/AgentSession.cs +++ b/src/LogAnalyzerAgent/Applications/AgentSession.cs @@ -79,22 +79,118 @@ 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 + { + _analyzer.ChangeDirectory(request.DirectoryPath); + 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 = CreateInternalErrorOperationStatus(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 = CreateInternalErrorOperationStatus(ex); + _logger.LogError(ex, "An error occurred while analyzing 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) || + result is null) + { + responses.Add(new GetAnalysisResultResponse + { + Status = new OperationStatusMessage + { + Success = false, + Code = AgentErrorCode.FileNotFound, + Message = $"File '{request.FileName}' was not found." + } + }); + } + else + { + 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 + { + Status = CreateNoErrorOperationStatus(), + Header = header + }); + + if (result.State == AnalysisState.Succeeded) + { + foreach (var entry in result.Entries) + { + cancellationToken.ThrowIfCancellationRequested(); + + responses.Add(new GetAnalysisResultResponse + { + Status = CreateNoErrorOperationStatus(), + LogEntry = GrpcTypeConverter.ConvertToGrpc(entry) + }); + } + } + } + } + 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..bb1116c 100644 --- a/src/LogAnalyzerAgent/Services/AgentService.cs +++ b/src/LogAnalyzerAgent/Services/AgentService.cs @@ -29,27 +29,33 @@ 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..c1da53d 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}"); + int separatorIndex = internalMessage.Exception.IndexOf(':'); // Find the first index of ":" + if(separatorIndex == -1) + { + throw new FormatException($"Invalid exception format: {internalMessage.Exception}"); + } + string exceptionName = internalMessage.Exception.Substring(0, separatorIndex); + string exceptionMessage = internalMessage.Exception.Substring(separatorIndex + 1).TrimStart(); + 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..9ce38c4 100644 --- a/src/RemoteCli/Program.cs +++ b/src/RemoteCli/Program.cs @@ -4,7 +4,6 @@ using LogAnalyzerRpc; using LogAnalyzerRpc.Protos; using LogParser.Visitors; -using Microsoft.Extensions.Logging; namespace RemoteCli { @@ -18,11 +17,23 @@ static async Task Main(string[] args) ?? Environment.GetEnvironmentVariable("LOG_ANALYZER_AGENT_ADDRESS") ?? "http://localhost:5000"; Console.WriteLine($"Connecting to agent at {address}..."); - using var channel = GrpcChannel.ForAddress(address); - var client = new LogAnalyzerAgentServiceClient(channel); - _ = await client.PingAsync(new Empty()); - await ChooseAction(client); + try + { + using var channel = GrpcChannel.ForAddress(address); + var client = new LogAnalyzerAgentServiceClient(channel); + _ = await client.PingAsync(new Empty()); + + await ChooseAction(client); + } + catch (RpcException ex) + { + PrintRpcError(ex); + } + catch (UriFormatException ex) + { + Console.WriteLine($"Invalid agent address: {ex.Message}"); + } } private static async Task InputDirectory(LogAnalyzerAgentServiceClient client) @@ -73,11 +84,7 @@ 6. Exit. { return; } - try - { - choice = int.Parse(choiceStr); - } - catch (Exception) + if (!int.TryParse(choiceStr, out choice)) { Console.WriteLine("Invalid input, please try again."); continue; @@ -90,58 +97,219 @@ 6. Exit. { 3, AnalyzeAll }, { 4, GetAnalysisResult } }; - switch (choice) + try { - case 1: - case 2: - case 3: - case 4: - await actions[choice](client); - break; - case 5: - var success = await InputDirectory(client); - if (!success) - { + switch (choice) + { + case 1: + case 2: + case 3: + case 4: + await actions[choice](client); + break; + case 5: + var success = await InputDirectory(client); + if (!success) + { + return; + } + break; + case 6: return; - } - break; - case 6: - return; - default: - Console.WriteLine("Invalid choice, please try again."); - break; + default: + Console.WriteLine("Invalid choice, please try again."); + break; + } + } + catch (RpcException ex) + { + PrintRpcError(ex); } } } private static async Task ShowLogFiles(LogAnalyzerAgentServiceClient client) { - throw new NotImplementedException("TODO: T3.2"); + var response = await client.GetLogFilesAsync(new Empty()); + if (!response.Status.Success) + { + PrintOperationError(response.Status); + return; + } + + Console.WriteLine($"[{string.Join(", ", response.FileNames)}]"); } private static int ReadDegreeOfParallelism() { - throw new NotImplementedException("TODO: T3.2"); + while (true) + { + Console.WriteLine("Please input the degree of parallelism (0 means auto):"); + Console.Write(">>> "); + Console.Out.Flush(); + + var input = Console.ReadLine(); + if (input is null) + { + return 0; + } + + if (int.TryParse(input, out var degreeOfParallelism) && degreeOfParallelism >= 0) + { + return degreeOfParallelism; + } + + Console.WriteLine("Invalid input, please try again."); + } } private static List ReadFileNames() { - throw new NotImplementedException("TODO: T3.2"); + Console.WriteLine("Please input file names to analyze, separated by commas:"); + Console.Write(">>> "); + Console.Out.Flush(); + + var input = Console.ReadLine(); + if (input is null) + { + return []; + } + + return [.. input.Split( + ',', + StringSplitOptions.TrimEntries | StringSplitOptions.RemoveEmptyEntries)]; } private static async Task AnalyzeFiles(LogAnalyzerAgentServiceClient client) { - throw new NotImplementedException("TODO: T3.2"); + var degreeOfParallelism = ReadDegreeOfParallelism(); + var fileNames = ReadFileNames(); + if (fileNames.Count == 0) + { + Console.WriteLine("No file names input."); + return; + } + + var request = new AnalyzeFilesRequest + { + DegreeOfParallelism = degreeOfParallelism, + }; + request.FileNames.AddRange(fileNames); + + var response = await client.AnalyzeFilesAsync(request); + if (!response.Status.Success) + { + PrintOperationError(response.Status); + return; + } + + Console.WriteLine($"Analysis finished: [{string.Join(", ", fileNames)}]"); } private static async Task AnalyzeAll(LogAnalyzerAgentServiceClient client) { - throw new NotImplementedException("TODO: T3.2"); + var request = new AnalyzeAllRequest + { + DegreeOfParallelism = ReadDegreeOfParallelism(), + }; + + var response = await client.AnalyzeAllAsync(request); + if (!response.Status.Success) + { + PrintOperationError(response.Status); + return; + } + + Console.WriteLine("Analysis finished."); } private static async Task GetAnalysisResult(LogAnalyzerAgentServiceClient client) { - throw new NotImplementedException("TODO: T3.2"); + Console.WriteLine("Please input the file name:"); + Console.Write(">>> "); + Console.Out.Flush(); + + var fileName = Console.ReadLine(); + if (fileName is null) + { + return; + } + + if (string.IsNullOrWhiteSpace(fileName)) + { + Console.WriteLine("File name cannot be empty."); + return; + } + + using var call = client.GetAnalysisResult(new GetAnalysisResultRequest + { + FileName = fileName, + }); + + var receivedResponse = false; + var dumper = new KeyValueVisitor(); + await foreach (var response in call.ResponseStream.ReadAllAsync()) + { + receivedResponse = true; + + if (!response.Status.Success) + { + PrintOperationError(response.Status); + return; + } + + switch (response.PayloadCase) + { + case GetAnalysisResultResponse.PayloadOneofCase.Header: + PrintAnalysisHeader(response.Header); + break; + case GetAnalysisResultResponse.PayloadOneofCase.LogEntry: + var entry = GrpcTypeConverter.ConvertFromGrpc(response.LogEntry); + var keyValuePairs = dumper.Dump(entry); + Console.WriteLine(string.Join(", ", + keyValuePairs.Select(pair => $"{pair.Key}: {pair.Value}"))); + break; + default: + Console.WriteLine("Error: The agent returned an invalid analysis result."); + return; + } + } + + if (!receivedResponse) + { + Console.WriteLine("Error: The agent returned no analysis result."); + } + } + + private static void PrintAnalysisHeader(AnalysisResultHeaderMessage header) + { + switch (header.State) + { + case AnalysisStateEnum.NotAnalyzed: + Console.WriteLine($"File '{header.FileName}' has not been analyzed."); + break; + case AnalysisStateEnum.Succeeded: + break; + case AnalysisStateEnum.Failed: + var errorMessage = header.HasErrorMessage + ? header.ErrorMessage + : "Unknown error."; + Console.WriteLine($"Analysis of file '{header.FileName}' failed: {errorMessage}"); + break; + default: + Console.WriteLine($"Error: Unknown analysis state '{header.State}'."); + break; + } + } + + private static void PrintOperationError(OperationStatusMessage status) + { + Console.WriteLine($"Error: {status.Code}: {status.Message}"); + } + + private static void PrintRpcError(RpcException exception) + { + Console.WriteLine($"RPC failed: {exception.StatusCode}: {exception.Status.Detail}"); } } }