mirror of
https://github.com/GreptimeTeam/greptimedb.git
synced 2026-05-20 23:10:37 +00:00
136 lines
15 KiB
HTML
136 lines
15 KiB
HTML
<!DOCTYPE html><html lang="en"><head><meta charset="utf-8"><meta name="viewport" content="width=device-width, initial-scale=1.0"><meta name="generator" content="rustdoc"><meta name="description" content="Source of the Rust file `src/common/procedure/src/watcher.rs`."><title>watcher.rs - source</title><script>if(window.location.protocol!=="file:")document.head.insertAdjacentHTML("beforeend","SourceSerif4-Regular-6b053e98.ttf.woff2,FiraSans-Italic-81dc35de.woff2,FiraSans-Regular-0fe48ade.woff2,FiraSans-MediumItalic-ccf7e434.woff2,FiraSans-Medium-e1aa3f0a.woff2,SourceCodePro-Regular-8badfe75.ttf.woff2,SourceCodePro-Semibold-aa29a496.ttf.woff2".split(",").map(f=>`<link rel="preload" as="font" type="font/woff2"href="../../static.files/${f}">`).join(""))</script><link rel="stylesheet" href="../../static.files/normalize-9960930a.css"><link rel="stylesheet" href="../../static.files/rustdoc-17e0aaed.css"><meta name="rustdoc-vars" data-root-path="../../" data-static-root-path="../../static.files/" data-current-crate="common_procedure" data-themes="" data-resource-suffix="" data-rustdoc-version="1.96.0-nightly (ac7f9ec7d 2026-03-20)" data-channel="nightly" data-search-js="search-63369b7b.js" data-stringdex-js="stringdex-2da4960a.js" data-settings-js="settings-170eb4bf.js" ><script src="../../static.files/storage-41dd4d93.js"></script><script defer src="../../static.files/src-script-813739b1.js"></script><script defer src="../../src-files.js"></script><script defer src="../../static.files/main-5013f961.js"></script><noscript><link rel="stylesheet" href="../../static.files/noscript-f7c3ffd8.css"></noscript><link rel="alternate icon" type="image/png" href="../../static.files/favicon-32x32-eab170b8.png"><link rel="icon" type="image/svg+xml" href="../../static.files/favicon-044be391.svg"></head><body class="rustdoc src"><a class="skip-main-content" href="#main-content">Skip to main content</a><!--[if lte IE 11]><div class="warning">This old browser is unsupported and will most likely display funky things.</div><![endif]--><nav class="sidebar"><div class="src-sidebar-title"><h2>Files</h2></div></nav><div class="sidebar-resizer" title="Drag to resize sidebar"></div><main><section id="main-content" class="content" tabindex="-1"><div class="main-heading"><h1><div class="sub-heading">common_procedure/</div>watcher.rs</h1><rustdoc-toolbar></rustdoc-toolbar></div><div class="example-wrap digits-3"><pre class="rust"><code><a href=#1 id=1 data-nosnippet>1</a><span class="comment">// Copyright 2023 Greptime Team
|
|
<a href=#2 id=2 data-nosnippet>2</a>//
|
|
<a href=#3 id=3 data-nosnippet>3</a>// Licensed under the Apache License, Version 2.0 (the "License");
|
|
<a href=#4 id=4 data-nosnippet>4</a>// you may not use this file except in compliance with the License.
|
|
<a href=#5 id=5 data-nosnippet>5</a>// You may obtain a copy of the License at
|
|
<a href=#6 id=6 data-nosnippet>6</a>//
|
|
<a href=#7 id=7 data-nosnippet>7</a>// http://www.apache.org/licenses/LICENSE-2.0
|
|
<a href=#8 id=8 data-nosnippet>8</a>//
|
|
<a href=#9 id=9 data-nosnippet>9</a>// Unless required by applicable law or agreed to in writing, software
|
|
<a href=#10 id=10 data-nosnippet>10</a>// distributed under the License is distributed on an "AS IS" BASIS,
|
|
<a href=#11 id=11 data-nosnippet>11</a>// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
|
<a href=#12 id=12 data-nosnippet>12</a>// See the License for the specific language governing permissions and
|
|
<a href=#13 id=13 data-nosnippet>13</a>// limitations under the License.
|
|
<a href=#14 id=14 data-nosnippet>14</a>
|
|
<a href=#15 id=15 data-nosnippet>15</a></span><span class="kw">use </span>common_telemetry::debug;
|
|
<a href=#16 id=16 data-nosnippet>16</a><span class="kw">use </span>snafu::ResultExt;
|
|
<a href=#17 id=17 data-nosnippet>17</a><span class="kw">use </span>tokio::sync::watch::Receiver;
|
|
<a href=#18 id=18 data-nosnippet>18</a>
|
|
<a href=#19 id=19 data-nosnippet>19</a><span class="kw">use </span><span class="kw">crate</span>::error::{ProcedureExecSnafu, <span class="prelude-ty">Result</span>, WaitWatcherSnafu};
|
|
<a href=#20 id=20 data-nosnippet>20</a><span class="kw">use </span><span class="kw">crate</span>::procedure::{Output, ProcedureState};
|
|
<a href=#21 id=21 data-nosnippet>21</a>
|
|
<a href=#22 id=22 data-nosnippet>22</a><span class="doccomment">/// Watcher to watch procedure state.
|
|
<a href=#23 id=23 data-nosnippet>23</a></span><span class="kw">pub type </span>Watcher = Receiver<ProcedureState>;
|
|
<a href=#24 id=24 data-nosnippet>24</a>
|
|
<a href=#25 id=25 data-nosnippet>25</a><span class="doccomment">/// Wait the [Watcher] until the [ProcedureState] is done.
|
|
<a href=#26 id=26 data-nosnippet>26</a></span><span class="kw">pub async fn </span>wait(watcher: <span class="kw-2">&mut </span>Watcher) -> <span class="prelude-ty">Result</span><<span class="prelude-ty">Option</span><Output>> {
|
|
<a href=#27 id=27 data-nosnippet>27</a> <span class="kw">loop </span>{
|
|
<a href=#28 id=28 data-nosnippet>28</a> watcher.changed().<span class="kw">await</span>.context(WaitWatcherSnafu)<span class="question-mark">?</span>;
|
|
<a href=#29 id=29 data-nosnippet>29</a> <span class="kw">match </span><span class="kw-2">&*</span>watcher.borrow() {
|
|
<a href=#30 id=30 data-nosnippet>30</a> ProcedureState::Running => (),
|
|
<a href=#31 id=31 data-nosnippet>31</a> ProcedureState::Done { output } => {
|
|
<a href=#32 id=32 data-nosnippet>32</a> <span class="kw">return </span><span class="prelude-val">Ok</span>(output.clone());
|
|
<a href=#33 id=33 data-nosnippet>33</a> }
|
|
<a href=#34 id=34 data-nosnippet>34</a> ProcedureState::Failed { error } => {
|
|
<a href=#35 id=35 data-nosnippet>35</a> <span class="kw">return </span><span class="prelude-val">Err</span>(error.clone()).context(ProcedureExecSnafu);
|
|
<a href=#36 id=36 data-nosnippet>36</a> }
|
|
<a href=#37 id=37 data-nosnippet>37</a> ProcedureState::Retrying { error } => {
|
|
<a href=#38 id=38 data-nosnippet>38</a> <span class="macro">debug!</span>(<span class="string">"retrying, source: {}"</span>, error)
|
|
<a href=#39 id=39 data-nosnippet>39</a> }
|
|
<a href=#40 id=40 data-nosnippet>40</a> ProcedureState::RollingBack { error } => {
|
|
<a href=#41 id=41 data-nosnippet>41</a> <span class="macro">debug!</span>(<span class="string">"rolling back, source: {:?}"</span>, error)
|
|
<a href=#42 id=42 data-nosnippet>42</a> }
|
|
<a href=#43 id=43 data-nosnippet>43</a> ProcedureState::PrepareRollback { error } => {
|
|
<a href=#44 id=44 data-nosnippet>44</a> <span class="macro">debug!</span>(<span class="string">"commit rollback, source: {}"</span>, error)
|
|
<a href=#45 id=45 data-nosnippet>45</a> }
|
|
<a href=#46 id=46 data-nosnippet>46</a> ProcedureState::Poisoned { error, .. } => {
|
|
<a href=#47 id=47 data-nosnippet>47</a> <span class="macro">debug!</span>(<span class="string">"poisoned, source: {}"</span>, error);
|
|
<a href=#48 id=48 data-nosnippet>48</a> <span class="kw">return </span><span class="prelude-val">Err</span>(error.clone()).context(ProcedureExecSnafu);
|
|
<a href=#49 id=49 data-nosnippet>49</a> }
|
|
<a href=#50 id=50 data-nosnippet>50</a> }
|
|
<a href=#51 id=51 data-nosnippet>51</a> }
|
|
<a href=#52 id=52 data-nosnippet>52</a>}
|
|
<a href=#53 id=53 data-nosnippet>53</a>
|
|
<a href=#54 id=54 data-nosnippet>54</a><span class="attr">#[cfg(test)]
|
|
<a href=#55 id=55 data-nosnippet>55</a></span><span class="kw">mod </span>tests {
|
|
<a href=#56 id=56 data-nosnippet>56</a>
|
|
<a href=#57 id=57 data-nosnippet>57</a> <span class="kw">use </span>std::sync::Arc;
|
|
<a href=#58 id=58 data-nosnippet>58</a> <span class="kw">use </span>std::time::Duration;
|
|
<a href=#59 id=59 data-nosnippet>59</a>
|
|
<a href=#60 id=60 data-nosnippet>60</a> <span class="kw">use </span>async_trait::async_trait;
|
|
<a href=#61 id=61 data-nosnippet>61</a> <span class="kw">use </span>common_error::mock::MockError;
|
|
<a href=#62 id=62 data-nosnippet>62</a> <span class="kw">use </span>common_error::status_code::StatusCode;
|
|
<a href=#63 id=63 data-nosnippet>63</a> <span class="kw">use </span>common_test_util::temp_dir::create_temp_dir;
|
|
<a href=#64 id=64 data-nosnippet>64</a>
|
|
<a href=#65 id=65 data-nosnippet>65</a> <span class="kw">use </span>super::<span class="kw-2">*</span>;
|
|
<a href=#66 id=66 data-nosnippet>66</a> <span class="kw">use </span><span class="kw">crate</span>::error::Error;
|
|
<a href=#67 id=67 data-nosnippet>67</a> <span class="kw">use </span><span class="kw">crate</span>::local::{LocalManager, ManagerConfig, test_util};
|
|
<a href=#68 id=68 data-nosnippet>68</a> <span class="kw">use </span><span class="kw">crate</span>::procedure::PoisonKeys;
|
|
<a href=#69 id=69 data-nosnippet>69</a> <span class="kw">use </span><span class="kw">crate</span>::store::state_store::ObjectStateStore;
|
|
<a href=#70 id=70 data-nosnippet>70</a> <span class="kw">use </span><span class="kw">crate</span>::test_util::InMemoryPoisonStore;
|
|
<a href=#71 id=71 data-nosnippet>71</a> <span class="kw">use </span>crate::{
|
|
<a href=#72 id=72 data-nosnippet>72</a> Context, LockKey, Procedure, ProcedureId, ProcedureManager, ProcedureWithId, Status,
|
|
<a href=#73 id=73 data-nosnippet>73</a> };
|
|
<a href=#74 id=74 data-nosnippet>74</a>
|
|
<a href=#75 id=75 data-nosnippet>75</a> <span class="attr">#[tokio::test]
|
|
<a href=#76 id=76 data-nosnippet>76</a> </span><span class="kw">async fn </span>test_success_after_retry() {
|
|
<a href=#77 id=77 data-nosnippet>77</a> <span class="kw">let </span>dir = create_temp_dir(<span class="string">"after_retry"</span>);
|
|
<a href=#78 id=78 data-nosnippet>78</a> <span class="kw">let </span>config = ManagerConfig {
|
|
<a href=#79 id=79 data-nosnippet>79</a> parent_path: <span class="string">"data/"</span>.to_string(),
|
|
<a href=#80 id=80 data-nosnippet>80</a> max_retry_times: <span class="number">3</span>,
|
|
<a href=#81 id=81 data-nosnippet>81</a> retry_delay: Duration::from_millis(<span class="number">500</span>),
|
|
<a href=#82 id=82 data-nosnippet>82</a> ..Default::default()
|
|
<a href=#83 id=83 data-nosnippet>83</a> };
|
|
<a href=#84 id=84 data-nosnippet>84</a> <span class="kw">let </span>state_store = Arc::new(ObjectStateStore::new(test_util::new_object_store(<span class="kw-2">&</span>dir)));
|
|
<a href=#85 id=85 data-nosnippet>85</a> <span class="kw">let </span>poison_manager = Arc::new(InMemoryPoisonStore::default());
|
|
<a href=#86 id=86 data-nosnippet>86</a> <span class="kw">let </span>manager = LocalManager::new(config, state_store, poison_manager, <span class="prelude-val">None</span>, <span class="prelude-val">None</span>);
|
|
<a href=#87 id=87 data-nosnippet>87</a> manager.start().<span class="kw">await</span>.unwrap();
|
|
<a href=#88 id=88 data-nosnippet>88</a>
|
|
<a href=#89 id=89 data-nosnippet>89</a> <span class="attr">#[derive(Debug)]
|
|
<a href=#90 id=90 data-nosnippet>90</a> </span><span class="kw">struct </span>MockProcedure {
|
|
<a href=#91 id=91 data-nosnippet>91</a> error: bool,
|
|
<a href=#92 id=92 data-nosnippet>92</a> }
|
|
<a href=#93 id=93 data-nosnippet>93</a>
|
|
<a href=#94 id=94 data-nosnippet>94</a> <span class="attr">#[async_trait]
|
|
<a href=#95 id=95 data-nosnippet>95</a> </span><span class="kw">impl </span>Procedure <span class="kw">for </span>MockProcedure {
|
|
<a href=#96 id=96 data-nosnippet>96</a> <span class="kw">fn </span>type_name(<span class="kw-2">&</span><span class="self">self</span>) -> <span class="kw-2">&</span>str {
|
|
<a href=#97 id=97 data-nosnippet>97</a> <span class="string">"MockProcedure"
|
|
<a href=#98 id=98 data-nosnippet>98</a> </span>}
|
|
<a href=#99 id=99 data-nosnippet>99</a>
|
|
<a href=#100 id=100 data-nosnippet>100</a> <span class="kw">async fn </span>execute(<span class="kw-2">&mut </span><span class="self">self</span>, _ctx: <span class="kw-2">&</span>Context) -> <span class="prelude-ty">Result</span><Status> {
|
|
<a href=#101 id=101 data-nosnippet>101</a> <span class="kw">if </span><span class="self">self</span>.error {
|
|
<a href=#102 id=102 data-nosnippet>102</a> <span class="self">self</span>.error = !<span class="self">self</span>.error;
|
|
<a href=#103 id=103 data-nosnippet>103</a> <span class="prelude-val">Err</span>(Error::retry_later(MockError::new(StatusCode::Internal)))
|
|
<a href=#104 id=104 data-nosnippet>104</a> } <span class="kw">else </span>{
|
|
<a href=#105 id=105 data-nosnippet>105</a> <span class="prelude-val">Ok</span>(Status::done_with_output(<span class="string">"hello"</span>))
|
|
<a href=#106 id=106 data-nosnippet>106</a> }
|
|
<a href=#107 id=107 data-nosnippet>107</a> }
|
|
<a href=#108 id=108 data-nosnippet>108</a>
|
|
<a href=#109 id=109 data-nosnippet>109</a> <span class="kw">fn </span>dump(<span class="kw-2">&</span><span class="self">self</span>) -> <span class="prelude-ty">Result</span><String> {
|
|
<a href=#110 id=110 data-nosnippet>110</a> <span class="prelude-val">Ok</span>(String::new())
|
|
<a href=#111 id=111 data-nosnippet>111</a> }
|
|
<a href=#112 id=112 data-nosnippet>112</a>
|
|
<a href=#113 id=113 data-nosnippet>113</a> <span class="kw">fn </span>lock_key(<span class="kw-2">&</span><span class="self">self</span>) -> LockKey {
|
|
<a href=#114 id=114 data-nosnippet>114</a> LockKey::single_exclusive(<span class="string">"test.submit"</span>)
|
|
<a href=#115 id=115 data-nosnippet>115</a> }
|
|
<a href=#116 id=116 data-nosnippet>116</a>
|
|
<a href=#117 id=117 data-nosnippet>117</a> <span class="kw">fn </span>poison_keys(<span class="kw-2">&</span><span class="self">self</span>) -> PoisonKeys {
|
|
<a href=#118 id=118 data-nosnippet>118</a> PoisonKeys::default()
|
|
<a href=#119 id=119 data-nosnippet>119</a> }
|
|
<a href=#120 id=120 data-nosnippet>120</a> }
|
|
<a href=#121 id=121 data-nosnippet>121</a>
|
|
<a href=#122 id=122 data-nosnippet>122</a> <span class="kw">let </span>procedure_id = ProcedureId::random();
|
|
<a href=#123 id=123 data-nosnippet>123</a> <span class="kw">let </span><span class="kw-2">mut </span>watcher = manager
|
|
<a href=#124 id=124 data-nosnippet>124</a> .submit(ProcedureWithId {
|
|
<a href=#125 id=125 data-nosnippet>125</a> id: procedure_id,
|
|
<a href=#126 id=126 data-nosnippet>126</a> procedure: Box::new(MockProcedure { error: <span class="bool-val">true </span>}),
|
|
<a href=#127 id=127 data-nosnippet>127</a> })
|
|
<a href=#128 id=128 data-nosnippet>128</a> .<span class="kw">await
|
|
<a href=#129 id=129 data-nosnippet>129</a> </span>.unwrap();
|
|
<a href=#130 id=130 data-nosnippet>130</a>
|
|
<a href=#131 id=131 data-nosnippet>131</a> <span class="kw">let </span>output = wait(<span class="kw-2">&mut </span>watcher).<span class="kw">await</span>.unwrap().unwrap();
|
|
<a href=#132 id=132 data-nosnippet>132</a> <span class="kw">let </span>output = output.downcast::<<span class="kw-2">&</span>str>().unwrap();
|
|
<a href=#133 id=133 data-nosnippet>133</a> <span class="macro">assert_eq!</span>(output.as_ref(), <span class="kw-2">&</span><span class="string">"hello"</span>);
|
|
<a href=#134 id=134 data-nosnippet>134</a> }
|
|
<a href=#135 id=135 data-nosnippet>135</a>}
|
|
</code></pre></div></section></main></body></html> |