From 5d9924026e7ef2db69c121280d70e330c45dd4ce Mon Sep 17 00:00:00 2001 From: Nanaloveyuki Date: Tue, 7 Jul 2026 11:35:16 +0800 Subject: [PATCH] =?UTF-8?q?=F0=9F=90=9B=20=E4=BF=AE=E5=A4=8D=E5=BC=82?= =?UTF-8?q?=E5=B8=B8=20Runtime=20Queue?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- src/runtime_file_controls.mbt | 111 +++++++++++++++------- src/runtime_logger.mbt | 173 +++++++++++++++++++++++++++++----- 2 files changed, 225 insertions(+), 59 deletions(-) diff --git a/src/runtime_file_controls.mbt b/src/runtime_file_controls.mbt index df980f9..2db0bb2 100644 --- a/src/runtime_file_controls.mbt +++ b/src/runtime_file_controls.mbt @@ -7,6 +7,34 @@ fn runtime_file_sink_internal(sink : RuntimeSink) -> FileSink? { } } +///| +fn queued_file_policy_mutation_allowed_internal( + sink : QueuedSink[FileSink], +) -> Bool { + sink.pending_count() == 0 +} + +///| +fn queued_file_flush_internal(sink : QueuedSink[FileSink]) -> Bool { + let before_write_failures = sink.sink.write_failures() + let before_flush_failures = sink.sink.flush_failures() + let before_rotation_failures = sink.sink.rotation_failures() + ignore(sink.flush()) + let flushed = sink.sink.flush() + sink.pending_count() == 0 && + flushed && + sink.sink.write_failures() == before_write_failures && + sink.sink.flush_failures() == before_flush_failures && + sink.sink.rotation_failures() == before_rotation_failures +} + +///| +fn queued_file_close_internal(sink : QueuedSink[FileSink]) -> Bool { + let flushed = queued_file_flush_internal(sink) + let closed = sink.sink.close() + flushed && closed +} + ///| pub fn RuntimeSink::file_available(self : RuntimeSink) -> Bool { match self { @@ -23,36 +51,29 @@ pub fn RuntimeSink::file_reopen( ) -> Bool { match self { File(sink) => sink.reopen(append~) - QueuedFile(sink) => sink.sink.reopen(append~) + QueuedFile(sink) => + if queued_file_policy_mutation_allowed_internal(sink) { + sink.sink.reopen(append~) + } else { + false + } _ => false } } ///| pub fn RuntimeSink::file_reopen_with_current_policy(self : RuntimeSink) -> Bool { - match self { - File(sink) => sink.reopen_with_current_policy() - QueuedFile(sink) => sink.sink.reopen_with_current_policy() - _ => false - } + self.file_reopen() } ///| pub fn RuntimeSink::file_reopen_append(self : RuntimeSink) -> Bool { - match self { - File(sink) => sink.reopen_append() - QueuedFile(sink) => sink.sink.reopen_append() - _ => false - } + self.file_reopen(append=Some(true)) } ///| pub fn RuntimeSink::file_reopen_truncate(self : RuntimeSink) -> Bool { - match self { - File(sink) => sink.reopen_truncate() - QueuedFile(sink) => sink.sink.reopen_truncate() - _ => false - } + self.file_reopen(append=Some(false)) } ///| @@ -75,8 +96,12 @@ pub fn RuntimeSink::file_set_append_mode( true } QueuedFile(sink) => { - sink.sink.set_append_mode(append) - true + if queued_file_policy_mutation_allowed_internal(sink) { + sink.sink.set_append_mode(append) + true + } else { + false + } } _ => false } @@ -130,8 +155,12 @@ pub fn RuntimeSink::file_set_auto_flush( true } QueuedFile(sink) => { - sink.sink.set_auto_flush(enabled) - true + if queued_file_policy_mutation_allowed_internal(sink) { + sink.sink.set_auto_flush(enabled) + true + } else { + false + } } _ => false } @@ -148,8 +177,12 @@ pub fn RuntimeSink::file_set_policy( true } QueuedFile(sink) => { - sink.sink.set_policy(policy) - true + if queued_file_policy_mutation_allowed_internal(sink) { + sink.sink.set_policy(policy) + true + } else { + false + } } _ => false } @@ -166,8 +199,12 @@ pub fn RuntimeSink::file_set_rotation( true } QueuedFile(sink) => { - sink.sink.set_rotation(rotation) - true + if queued_file_policy_mutation_allowed_internal(sink) { + sink.sink.set_rotation(rotation) + true + } else { + false + } } _ => false } @@ -181,8 +218,12 @@ pub fn RuntimeSink::file_clear_rotation(self : RuntimeSink) -> Bool { true } QueuedFile(sink) => { - sink.sink.clear_rotation() - true + if queued_file_policy_mutation_allowed_internal(sink) { + sink.sink.clear_rotation() + true + } else { + false + } } _ => false } @@ -192,10 +233,7 @@ pub fn RuntimeSink::file_clear_rotation(self : RuntimeSink) -> Bool { pub fn RuntimeSink::file_flush(self : RuntimeSink) -> Bool { match self { File(sink) => sink.flush() - QueuedFile(sink) => { - ignore(sink.flush()) - sink.sink.flush() - } + QueuedFile(sink) => queued_file_flush_internal(sink) _ => false } } @@ -204,10 +242,7 @@ pub fn RuntimeSink::file_flush(self : RuntimeSink) -> Bool { pub fn RuntimeSink::file_close(self : RuntimeSink) -> Bool { match self { File(sink) => sink.close() - QueuedFile(sink) => { - ignore(sink.flush()) - sink.sink.close() - } + QueuedFile(sink) => queued_file_close_internal(sink) _ => false } } @@ -271,8 +306,12 @@ pub fn RuntimeSink::file_reset_policy(self : RuntimeSink) -> Bool { true } QueuedFile(sink) => { - sink.sink.reset_policy() - true + if queued_file_policy_mutation_allowed_internal(sink) { + sink.sink.reset_policy() + true + } else { + false + } } _ => false } diff --git a/src/runtime_logger.mbt b/src/runtime_logger.mbt index a7f988b..740a487 100644 --- a/src/runtime_logger.mbt +++ b/src/runtime_logger.mbt @@ -13,6 +13,14 @@ pub(all) enum RuntimeSink { ///| pub type RuntimeFileState = @utils.RuntimeFileState +///| +pub struct RuntimeSinkProgress { + queue_advanced_count : Int + file_flush_step_count : Int + queue_backed : Bool + file_backed : Bool +} + ///| pub fn file_sink_policy_to_json( policy : FileSinkPolicy, @@ -71,31 +79,135 @@ pub impl Sink for RuntimeSink with fn write(self, rec) { } ///| -pub fn RuntimeSink::flush(self : RuntimeSink) -> Int { +fn[S : Sink] queued_sink_close_internal(sink : QueuedSink[S]) -> Bool { + ignore(sink.flush()) + sink.pending_count() == 0 +} + +///| +fn runtime_sink_progress_compat_count(progress : RuntimeSinkProgress) -> Int { + progress.queue_advanced_count + progress.file_flush_step_count +} + +///| +pub fn RuntimeSink::flush_progress(self : RuntimeSink) -> RuntimeSinkProgress { match self { - Console(_) => 0 - JsonConsole(_) => 0 - TextConsole(_) => 0 - File(sink) => if sink.flush() { 1 } else { 0 } - QueuedConsole(sink) => sink.flush() - QueuedJsonConsole(sink) => sink.flush() - QueuedTextConsole(sink) => sink.flush() - QueuedFile(sink) => sink.flush() + Console(_) => { + queue_advanced_count: 0, + file_flush_step_count: 0, + queue_backed: false, + file_backed: false, + } + JsonConsole(_) => { + queue_advanced_count: 0, + file_flush_step_count: 0, + queue_backed: false, + file_backed: false, + } + TextConsole(_) => { + queue_advanced_count: 0, + file_flush_step_count: 0, + queue_backed: false, + file_backed: false, + } + File(sink) => { + queue_advanced_count: 0, + file_flush_step_count: if sink.flush() { 1 } else { 0 }, + queue_backed: false, + file_backed: true, + } + QueuedConsole(sink) => { + queue_advanced_count: sink.flush(), + file_flush_step_count: 0, + queue_backed: true, + file_backed: false, + } + QueuedJsonConsole(sink) => { + queue_advanced_count: sink.flush(), + file_flush_step_count: 0, + queue_backed: true, + file_backed: false, + } + QueuedTextConsole(sink) => { + queue_advanced_count: sink.flush(), + file_flush_step_count: 0, + queue_backed: true, + file_backed: false, + } + QueuedFile(sink) => { + queue_advanced_count: sink.flush(), + file_flush_step_count: 0, + queue_backed: true, + file_backed: true, + } + } +} + +///| +pub fn RuntimeSink::flush(self : RuntimeSink) -> Int { + runtime_sink_progress_compat_count(self.flush_progress()) +} + +///| +pub fn RuntimeSink::drain_progress( + self : RuntimeSink, + max_items? : Int = -1, +) -> RuntimeSinkProgress { + match self { + Console(_) => { + queue_advanced_count: 0, + file_flush_step_count: 0, + queue_backed: false, + file_backed: false, + } + JsonConsole(_) => { + queue_advanced_count: 0, + file_flush_step_count: 0, + queue_backed: false, + file_backed: false, + } + TextConsole(_) => { + queue_advanced_count: 0, + file_flush_step_count: 0, + queue_backed: false, + file_backed: false, + } + File(sink) => { + queue_advanced_count: 0, + file_flush_step_count: if sink.flush() { 1 } else { 0 }, + queue_backed: false, + file_backed: true, + } + QueuedConsole(sink) => { + queue_advanced_count: sink.drain(max_items~), + file_flush_step_count: 0, + queue_backed: true, + file_backed: false, + } + QueuedJsonConsole(sink) => { + queue_advanced_count: sink.drain(max_items~), + file_flush_step_count: 0, + queue_backed: true, + file_backed: false, + } + QueuedTextConsole(sink) => { + queue_advanced_count: sink.drain(max_items~), + file_flush_step_count: 0, + queue_backed: true, + file_backed: false, + } + QueuedFile(sink) => { + queue_advanced_count: sink.drain(max_items~), + file_flush_step_count: 0, + queue_backed: true, + file_backed: true, + } } } ///| pub fn RuntimeSink::drain(self : RuntimeSink, max_items? : Int = -1) -> Int { - match self { - Console(_) => 0 - JsonConsole(_) => 0 - TextConsole(_) => 0 - File(sink) => if sink.flush() { 1 } else { 0 } - QueuedConsole(sink) => sink.drain(max_items~) - QueuedJsonConsole(sink) => sink.drain(max_items~) - QueuedTextConsole(sink) => sink.drain(max_items~) - QueuedFile(sink) => sink.drain(max_items~) - } + runtime_sink_progress_compat_count(self.drain_progress(max_items~)) } ///| @@ -105,10 +217,10 @@ pub fn RuntimeSink::close(self : RuntimeSink) -> Bool { JsonConsole(_) => true TextConsole(_) => true File(sink) => sink.close() - QueuedConsole(_) => true - QueuedJsonConsole(_) => true - QueuedTextConsole(_) => true - QueuedFile(sink) => sink.sink.close() + QueuedConsole(sink) => queued_sink_close_internal(sink) + QueuedJsonConsole(sink) => queued_sink_close_internal(sink) + QueuedTextConsole(sink) => queued_sink_close_internal(sink) + QueuedFile(sink) => queued_file_close_internal(sink) } } @@ -143,11 +255,26 @@ pub fn RuntimeSink::dropped_count(self : RuntimeSink) -> Int { ///| pub type ConfiguredLogger = Logger[RuntimeSink] +///| +pub fn ConfiguredLogger::flush_progress( + self : ConfiguredLogger, +) -> RuntimeSinkProgress { + self.sink.flush_progress() +} + ///| pub fn ConfiguredLogger::flush(self : ConfiguredLogger) -> Int { self.sink.flush() } +///| +pub fn ConfiguredLogger::drain_progress( + self : ConfiguredLogger, + max_items? : Int = -1, +) -> RuntimeSinkProgress { + self.sink.drain_progress(max_items~) +} + ///| pub fn ConfiguredLogger::drain( self : ConfiguredLogger,