file_sink_test.go 22 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818
  1. // Copyright 2023 EMQ Technologies Co., Ltd.
  2. //
  3. // Licensed under the Apache License, Version 2.0 (the "License");
  4. // you may not use this file except in compliance with the License.
  5. // You may obtain a copy of the License at
  6. //
  7. // http://www.apache.org/licenses/LICENSE-2.0
  8. //
  9. // Unless required by applicable law or agreed to in writing, software
  10. // distributed under the License is distributed on an "AS IS" BASIS,
  11. // WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
  12. // See the License for the specific language governing permissions and
  13. // limitations under the License.
  14. package file
  15. import (
  16. "fmt"
  17. "os"
  18. "path/filepath"
  19. "reflect"
  20. "strconv"
  21. "testing"
  22. "time"
  23. "github.com/lf-edge/ekuiper/internal/compressor"
  24. "github.com/lf-edge/ekuiper/internal/conf"
  25. "github.com/lf-edge/ekuiper/internal/topo/context"
  26. "github.com/lf-edge/ekuiper/internal/topo/topotest/mockclock"
  27. "github.com/lf-edge/ekuiper/internal/topo/transform"
  28. "github.com/lf-edge/ekuiper/pkg/message"
  29. )
  30. // Unit test for Configure function
  31. func TestConfigure(t *testing.T) {
  32. props := map[string]interface{}{
  33. "interval": 500,
  34. "path": "test",
  35. }
  36. m := File().(*fileSink)
  37. err := m.Configure(props)
  38. if err != nil {
  39. t.Errorf("Configure() error = %v, wantErr nil", err)
  40. }
  41. if m.c.Path != "test" {
  42. t.Errorf("Configure() Path = %v, want test", m.c.Path)
  43. }
  44. err = m.Configure(map[string]interface{}{"interval": 500, "path": ""})
  45. if err == nil {
  46. t.Errorf("Configure() error = %v, wantErr not nil", err)
  47. }
  48. err = m.Configure(map[string]interface{}{"fileType": "csv2"})
  49. if err == nil {
  50. t.Errorf("Configure() error = %v, wantErr not nil", err)
  51. }
  52. err = m.Configure(map[string]interface{}{
  53. "interval": 500,
  54. "path": "test",
  55. "fileType": "csv",
  56. })
  57. if err == nil {
  58. t.Errorf("Configure() error = %v, wantErr not nil", err)
  59. }
  60. err = m.Configure(map[string]interface{}{"interval": 60, "path": "test", "checkInterval": -1})
  61. if err == nil {
  62. t.Errorf("Configure() error = %v, wantErr not nil", err)
  63. }
  64. err = m.Configure(map[string]interface{}{"rollingInterval": -1})
  65. if err == nil {
  66. t.Errorf("Configure() error = %v, wantErr not nil", err)
  67. }
  68. err = m.Configure(map[string]interface{}{"rollingCount": -1})
  69. if err == nil {
  70. t.Errorf("Configure() error = %v, wantErr not nil", err)
  71. }
  72. err = m.Configure(map[string]interface{}{"rollingCount": 0, "rollingInterval": 0})
  73. if err == nil {
  74. t.Errorf("Configure() error = %v, wantErr not nil", err)
  75. }
  76. err = m.Configure(map[string]interface{}{"RollingNamePattern": "test"})
  77. if err == nil {
  78. t.Errorf("Configure() error = %v, wantErr not nil", err)
  79. }
  80. err = m.Configure(map[string]interface{}{"RollingNamePattern": 0})
  81. if err == nil {
  82. t.Errorf("Configure() error = %v, wantErr not nil", err)
  83. }
  84. for k := range compressionTypes {
  85. err = m.Configure(map[string]interface{}{
  86. "interval": 500,
  87. "path": "test",
  88. "compression": k,
  89. "rollingNamePattern": "suffix",
  90. })
  91. if err != nil {
  92. t.Errorf("Configure() error = %v, wantErr nil", err)
  93. }
  94. if m.c.Compression != k {
  95. t.Errorf("Configure() Compression = %v, want %v", m.c.Compression, k)
  96. }
  97. }
  98. err = m.Configure(map[string]interface{}{
  99. "interval": 500,
  100. "path": "test",
  101. "compression": "",
  102. "rollingNamePattern": "suffix",
  103. })
  104. if err != nil {
  105. t.Errorf("Configure() error = %v, wantErr nil", err)
  106. }
  107. if m.c.Compression != "" {
  108. t.Errorf("Configure() Compression = %v, want %v", m.c.Compression, "")
  109. }
  110. err = m.Configure(map[string]interface{}{
  111. "interval": 500,
  112. "path": "test",
  113. "compression": "not_exist_algorithm",
  114. })
  115. if err == nil {
  116. t.Errorf("Configure() error = %v, wantErr not nil", err)
  117. }
  118. }
  119. func TestFileSink_Configure(t *testing.T) {
  120. var (
  121. defaultCheckInterval = (5 * time.Minute).Milliseconds()
  122. int500 = 500
  123. int64_500 = int64(int500)
  124. )
  125. tests := []struct {
  126. name string
  127. c *sinkConf
  128. p map[string]interface{}
  129. }{
  130. {
  131. name: "default configurations",
  132. c: &sinkConf{
  133. CheckInterval: &defaultCheckInterval,
  134. Path: "cache",
  135. FileType: LINES_TYPE,
  136. RollingCount: 1000000,
  137. },
  138. p: map[string]interface{}{},
  139. },
  140. {
  141. name: "new props",
  142. c: &sinkConf{
  143. CheckInterval: &int64_500,
  144. Path: "test",
  145. FileType: CSV_TYPE,
  146. Format: message.FormatDelimited,
  147. Delimiter: ",",
  148. RollingCount: 1000000,
  149. RollingNamePattern: "none",
  150. },
  151. p: map[string]interface{}{
  152. "checkInterval": 500,
  153. "path": "test",
  154. "fileType": "csv",
  155. "format": message.FormatDelimited,
  156. "rollingNamePattern": "none",
  157. },
  158. },
  159. { // only set rolling interval
  160. name: "rolling",
  161. c: &sinkConf{
  162. CheckInterval: &defaultCheckInterval,
  163. Path: "cache",
  164. FileType: LINES_TYPE,
  165. RollingInterval: 500,
  166. RollingCount: 0,
  167. },
  168. p: map[string]interface{}{
  169. "rollingInterval": 500,
  170. "rollingCount": 0,
  171. },
  172. },
  173. {
  174. name: "fields",
  175. c: &sinkConf{
  176. CheckInterval: &defaultCheckInterval,
  177. Path: "cache",
  178. FileType: LINES_TYPE,
  179. RollingInterval: 500,
  180. RollingCount: 0,
  181. Fields: []string{"c", "a", "b"},
  182. },
  183. p: map[string]interface{}{
  184. "rollingInterval": 500,
  185. "rollingCount": 0,
  186. "fields": []string{"c", "a", "b"},
  187. },
  188. },
  189. }
  190. for _, tt := range tests {
  191. t.Run(tt.name, func(t *testing.T) {
  192. m := &fileSink{}
  193. if err := m.Configure(tt.p); err != nil {
  194. t.Errorf("fileSink.Configure() error = %v", err)
  195. return
  196. }
  197. if !reflect.DeepEqual(m.c, tt.c) {
  198. t.Errorf("fileSink.Configure() = %v, want %v", m.c, tt.c)
  199. }
  200. })
  201. }
  202. }
  203. // Test single file writing and flush by close
  204. func TestFileSink_Collect(t *testing.T) {
  205. tests := []struct {
  206. name string
  207. ft FileType
  208. fname string
  209. content []byte
  210. compress string
  211. }{
  212. {
  213. name: "lines",
  214. ft: LINES_TYPE,
  215. fname: "test_lines",
  216. content: []byte("{\"key\":\"value1\"}\n{\"key\":\"value2\"}"),
  217. },
  218. {
  219. name: "json",
  220. ft: JSON_TYPE,
  221. fname: "test_json",
  222. content: []byte(`[{"key":"value1"},{"key":"value2"}]`),
  223. },
  224. {
  225. name: "csv",
  226. ft: CSV_TYPE,
  227. fname: "test_csv",
  228. content: []byte("key\n{\"key\":\"value1\"}\n{\"key\":\"value2\"}"),
  229. },
  230. {
  231. name: "lines",
  232. ft: LINES_TYPE,
  233. fname: "test_lines",
  234. content: []byte("{\"key\":\"value1\"}\n{\"key\":\"value2\"}"),
  235. compress: GZIP,
  236. },
  237. {
  238. name: "json",
  239. ft: JSON_TYPE,
  240. fname: "test_json",
  241. content: []byte(`[{"key":"value1"},{"key":"value2"}]`),
  242. compress: GZIP,
  243. },
  244. {
  245. name: "csv",
  246. ft: CSV_TYPE,
  247. fname: "test_csv",
  248. content: []byte("key\n{\"key\":\"value1\"}\n{\"key\":\"value2\"}"),
  249. compress: GZIP,
  250. },
  251. {
  252. name: "lines",
  253. ft: LINES_TYPE,
  254. fname: "test_lines",
  255. content: []byte("{\"key\":\"value1\"}\n{\"key\":\"value2\"}"),
  256. compress: ZSTD,
  257. },
  258. {
  259. name: "json",
  260. ft: JSON_TYPE,
  261. fname: "test_json",
  262. content: []byte(`[{"key":"value1"},{"key":"value2"}]`),
  263. compress: ZSTD,
  264. },
  265. {
  266. name: "csv",
  267. ft: CSV_TYPE,
  268. fname: "test_csv",
  269. content: []byte("key\n{\"key\":\"value1\"}\n{\"key\":\"value2\"}"),
  270. compress: ZSTD,
  271. },
  272. }
  273. // Create a stream context for testing
  274. contextLogger := conf.Log.WithField("rule", "test2")
  275. ctx := context.WithValue(context.Background(), context.LoggerKey, contextLogger)
  276. tf, _ := transform.GenTransform("", "json", "", "", "", []string{})
  277. vCtx := context.WithValue(ctx, context.TransKey, tf)
  278. for _, tt := range tests {
  279. t.Run(tt.name, func(t *testing.T) {
  280. // Create a temporary file for testing
  281. tmpfile, err := os.CreateTemp("", tt.fname)
  282. if err != nil {
  283. t.Fatal(err)
  284. }
  285. defer os.Remove(tmpfile.Name())
  286. // Create a file sink with the temporary file path
  287. sink := &fileSink{}
  288. f := message.FormatJson
  289. if tt.ft == CSV_TYPE {
  290. f = message.FormatDelimited
  291. }
  292. err = sink.Configure(map[string]interface{}{
  293. "path": tmpfile.Name(),
  294. "fileType": tt.ft,
  295. "hasHeader": true,
  296. "format": f,
  297. "rollingNamePattern": "none",
  298. "compression": tt.compress,
  299. "fields": []string{"key"},
  300. })
  301. if err != nil {
  302. t.Fatal(err)
  303. }
  304. err = sink.Open(ctx)
  305. if err != nil {
  306. t.Fatal(err)
  307. }
  308. // Test collecting a map item
  309. m := map[string]interface{}{"key": "value1"}
  310. if err := sink.Collect(vCtx, m); err != nil {
  311. t.Errorf("unexpected error: %s", err)
  312. }
  313. // Test collecting another map item
  314. m = map[string]interface{}{"key": "value2"}
  315. if err := sink.Collect(ctx, m); err != nil {
  316. t.Errorf("unexpected error: %s", err)
  317. }
  318. if err = sink.Close(ctx); err != nil {
  319. t.Errorf("unexpected close error: %s", err)
  320. }
  321. // Read the contents of the temporary file and check if they match the collected items
  322. contents, err := os.ReadFile(tmpfile.Name())
  323. if err != nil {
  324. t.Fatal(err)
  325. }
  326. if tt.compress != "" {
  327. decompressor, _ := compressor.GetDecompressor(tt.compress)
  328. decompress, err := decompressor.Decompress(contents)
  329. if err != nil {
  330. t.Errorf("%v", err)
  331. }
  332. if !reflect.DeepEqual(decompress, tt.content) {
  333. t.Errorf("\nexpected\t %q \nbut got\t\t %q", tt.content, string(contents))
  334. }
  335. } else {
  336. if !reflect.DeepEqual(contents, tt.content) {
  337. t.Errorf("\nexpected\t %q \nbut got\t\t %q", tt.content, string(contents))
  338. }
  339. }
  340. })
  341. }
  342. }
  343. // Test file collect by fields
  344. func TestFileSinkFields_Collect(t *testing.T) {
  345. tests := []struct {
  346. name string
  347. ft FileType
  348. fname string
  349. format string
  350. delimiter string
  351. fields []string
  352. content []byte
  353. }{
  354. {
  355. name: "test1",
  356. ft: CSV_TYPE,
  357. fname: "test_csv",
  358. format: "delimited",
  359. delimiter: ",",
  360. fields: []string{"temperature", "humidity"},
  361. content: []byte("temperature,humidity\n31.2,40"),
  362. },
  363. {
  364. name: "test2",
  365. ft: CSV_TYPE,
  366. fname: "test_csv",
  367. format: "delimited",
  368. delimiter: ",",
  369. content: []byte("humidity,temperature\n40,31.2"),
  370. },
  371. }
  372. // Create a stream context for testing
  373. contextLogger := conf.Log.WithField("rule", "testFields")
  374. ctx := context.WithValue(context.Background(), context.LoggerKey, contextLogger)
  375. for _, tt := range tests {
  376. t.Run(tt.name, func(t *testing.T) {
  377. tf, _ := transform.GenTransform("", tt.format, "", tt.delimiter, "", tt.fields)
  378. vCtx := context.WithValue(ctx, context.TransKey, tf)
  379. // Create a temporary file for testing
  380. tmpfile, err := os.CreateTemp("", tt.fname)
  381. if err != nil {
  382. t.Fatal(err)
  383. }
  384. defer os.Remove(tmpfile.Name())
  385. // Create a file sink with the temporary file path
  386. sink := &fileSink{}
  387. err = sink.Configure(map[string]interface{}{
  388. "path": tmpfile.Name(),
  389. "fileType": tt.ft,
  390. "hasHeader": true,
  391. "format": tt.format,
  392. "rollingNamePattern": "none",
  393. "fields": tt.fields,
  394. })
  395. if err != nil {
  396. t.Fatal(err)
  397. }
  398. err = sink.Open(ctx)
  399. if err != nil {
  400. t.Fatal(err)
  401. }
  402. // Test collecting a map item
  403. m := map[string]interface{}{"humidity": 40, "temperature": 31.2}
  404. if err := sink.Collect(vCtx, m); err != nil {
  405. t.Errorf("unexpected error: %s", err)
  406. }
  407. if err = sink.Close(ctx); err != nil {
  408. t.Errorf("unexpected close error: %s", err)
  409. }
  410. // Read the contents of the temporary file and check if they match the collected items
  411. contents, err := os.ReadFile(tmpfile.Name())
  412. if err != nil {
  413. t.Fatal(err)
  414. }
  415. if !reflect.DeepEqual(contents, tt.content) {
  416. t.Errorf("\nexpected\t %q \nbut got\t\t %q", tt.content, string(contents))
  417. }
  418. })
  419. }
  420. }
  421. // Test file rolling by time
  422. func TestFileSinkRolling_Collect(t *testing.T) {
  423. // Remove existing files
  424. err := filepath.Walk(".", func(path string, info os.FileInfo, err error) error {
  425. if err != nil {
  426. return err
  427. }
  428. if filepath.Ext(path) == ".log" {
  429. fmt.Println("Deleting file:", path)
  430. return os.Remove(path)
  431. }
  432. return nil
  433. })
  434. if err != nil {
  435. t.Fatal(err)
  436. }
  437. conf.IsTesting = true
  438. tests := []struct {
  439. name string
  440. ft FileType
  441. fname string
  442. contents [2][]byte
  443. compress string
  444. }{
  445. {
  446. name: "lines",
  447. ft: LINES_TYPE,
  448. fname: "test_lines.log",
  449. contents: [2][]byte{
  450. []byte("{\"key\":\"value0\",\"ts\":460}\n{\"key\":\"value1\",\"ts\":910}\n{\"key\":\"value2\",\"ts\":1360}"),
  451. []byte("{\"key\":\"value3\",\"ts\":1810}\n{\"key\":\"value4\",\"ts\":2260}"),
  452. },
  453. },
  454. {
  455. name: "json",
  456. ft: JSON_TYPE,
  457. fname: "test_json.log",
  458. contents: [2][]byte{
  459. []byte("[{\"key\":\"value0\",\"ts\":460},{\"key\":\"value1\",\"ts\":910},{\"key\":\"value2\",\"ts\":1360}]"),
  460. []byte("[{\"key\":\"value3\",\"ts\":1810},{\"key\":\"value4\",\"ts\":2260}]"),
  461. },
  462. },
  463. {
  464. name: "lines",
  465. ft: LINES_TYPE,
  466. fname: "test_lines_gzip.log",
  467. contents: [2][]byte{
  468. []byte("{\"key\":\"value0\",\"ts\":460}\n{\"key\":\"value1\",\"ts\":910}\n{\"key\":\"value2\",\"ts\":1360}"),
  469. []byte("{\"key\":\"value3\",\"ts\":1810}\n{\"key\":\"value4\",\"ts\":2260}"),
  470. },
  471. compress: GZIP,
  472. },
  473. {
  474. name: "json",
  475. ft: JSON_TYPE,
  476. fname: "test_json_gzip.log",
  477. contents: [2][]byte{
  478. []byte("[{\"key\":\"value0\",\"ts\":460},{\"key\":\"value1\",\"ts\":910},{\"key\":\"value2\",\"ts\":1360}]"),
  479. []byte("[{\"key\":\"value3\",\"ts\":1810},{\"key\":\"value4\",\"ts\":2260}]"),
  480. },
  481. compress: GZIP,
  482. },
  483. {
  484. name: "lines",
  485. ft: LINES_TYPE,
  486. fname: "test_lines_zstd.log",
  487. contents: [2][]byte{
  488. []byte("{\"key\":\"value0\",\"ts\":460}\n{\"key\":\"value1\",\"ts\":910}\n{\"key\":\"value2\",\"ts\":1360}"),
  489. []byte("{\"key\":\"value3\",\"ts\":1810}\n{\"key\":\"value4\",\"ts\":2260}"),
  490. },
  491. compress: ZSTD,
  492. },
  493. {
  494. name: "json",
  495. ft: JSON_TYPE,
  496. fname: "test_json_zstd.log",
  497. contents: [2][]byte{
  498. []byte("[{\"key\":\"value0\",\"ts\":460},{\"key\":\"value1\",\"ts\":910},{\"key\":\"value2\",\"ts\":1360}]"),
  499. []byte("[{\"key\":\"value3\",\"ts\":1810},{\"key\":\"value4\",\"ts\":2260}]"),
  500. },
  501. compress: ZSTD,
  502. },
  503. }
  504. // Create a stream context for testing
  505. contextLogger := conf.Log.WithField("rule", "testRolling")
  506. ctx := context.WithValue(context.Background(), context.LoggerKey, contextLogger)
  507. tf, _ := transform.GenTransform("", "json", "", "", "", []string{})
  508. vCtx := context.WithValue(ctx, context.TransKey, tf)
  509. for _, tt := range tests {
  510. t.Run(tt.name, func(t *testing.T) {
  511. // Create a file sink with the temporary file path
  512. sink := &fileSink{}
  513. err := sink.Configure(map[string]interface{}{
  514. "path": tt.fname,
  515. "fileType": tt.ft,
  516. "rollingInterval": 1000,
  517. "checkInterval": 500,
  518. "rollingCount": 0,
  519. "rollingNamePattern": "suffix",
  520. "compression": tt.compress,
  521. })
  522. if err != nil {
  523. t.Fatal(err)
  524. }
  525. mockclock.ResetClock(10)
  526. err = sink.Open(ctx)
  527. if err != nil {
  528. t.Fatal(err)
  529. }
  530. c := mockclock.GetMockClock()
  531. for i := 0; i < 5; i++ {
  532. c.Add(450 * time.Millisecond)
  533. m := map[string]interface{}{"key": "value" + strconv.Itoa(i), "ts": c.Now().UnixMilli()}
  534. if err := sink.Collect(vCtx, m); err != nil {
  535. t.Errorf("unexpected error: %s", err)
  536. }
  537. }
  538. c.After(2000 * time.Millisecond)
  539. if err = sink.Close(ctx); err != nil {
  540. t.Errorf("unexpected close error: %s", err)
  541. }
  542. // Should write to 2 files
  543. for i := 0; i < 2; i++ {
  544. // Read the contents of the temporary file and check if they match the collected items
  545. var fn string
  546. if tt.compress != "" {
  547. fn = fmt.Sprintf("test_%s_%s-%d.log", tt.ft, tt.compress, 460+1350*i)
  548. } else {
  549. fn = fmt.Sprintf("test_%s-%d.log", tt.ft, 460+1350*i)
  550. }
  551. var contents []byte
  552. contents, err := os.ReadFile(fn)
  553. if err != nil {
  554. t.Fatal(err)
  555. }
  556. if tt.compress != "" {
  557. decompressor, _ := compressor.GetDecompressor(tt.compress)
  558. contents, err = decompressor.Decompress(contents)
  559. if err != nil {
  560. t.Errorf("%v", err)
  561. }
  562. }
  563. if !reflect.DeepEqual(contents, tt.contents[i]) {
  564. t.Errorf("\nexpected\t %q \nbut got\t\t %q", tt.contents[i], string(contents))
  565. }
  566. }
  567. })
  568. }
  569. }
  570. // Test file rolling by count
  571. func TestFileSinkRollingCount_Collect(t *testing.T) {
  572. // Remove existing files
  573. err := filepath.Walk(".", func(path string, info os.FileInfo, err error) error {
  574. if err != nil {
  575. return err
  576. }
  577. if filepath.Ext(path) == ".dd" {
  578. fmt.Println("Deleting file:", path)
  579. return os.Remove(path)
  580. }
  581. return nil
  582. })
  583. if err != nil {
  584. t.Fatal(err)
  585. }
  586. conf.IsTesting = true
  587. tests := []struct {
  588. name string
  589. ft FileType
  590. fname string
  591. contents [3][]byte
  592. compress string
  593. }{
  594. {
  595. name: "csv",
  596. ft: CSV_TYPE,
  597. fname: "test_csv_{{.ts}}.dd",
  598. contents: [3][]byte{
  599. []byte("key,ts\nvalue0,460"),
  600. []byte("key,ts\nvalue1,910"),
  601. []byte("key,ts\nvalue2,1360"),
  602. },
  603. },
  604. {
  605. name: "csv",
  606. ft: CSV_TYPE,
  607. fname: "test_csv_gzip_{{.ts}}.dd",
  608. contents: [3][]byte{
  609. []byte("key,ts\nvalue0,460"),
  610. []byte("key,ts\nvalue1,910"),
  611. []byte("key,ts\nvalue2,1360"),
  612. },
  613. compress: GZIP,
  614. },
  615. {
  616. name: "csv",
  617. ft: CSV_TYPE,
  618. fname: "test_csv_zstd_{{.ts}}.dd",
  619. contents: [3][]byte{
  620. []byte("key,ts\nvalue0,460"),
  621. []byte("key,ts\nvalue1,910"),
  622. []byte("key,ts\nvalue2,1360"),
  623. },
  624. compress: ZSTD,
  625. },
  626. }
  627. // Create a stream context for testing
  628. contextLogger := conf.Log.WithField("rule", "testRollingCount")
  629. ctx := context.WithValue(context.Background(), context.LoggerKey, contextLogger)
  630. tf, _ := transform.GenTransform("", "delimited", "", ",", "", []string{})
  631. vCtx := context.WithValue(ctx, context.TransKey, tf)
  632. for _, tt := range tests {
  633. t.Run(tt.name, func(t *testing.T) {
  634. // Create a file sink with the temporary file path
  635. sink := &fileSink{}
  636. err := sink.Configure(map[string]interface{}{
  637. "path": tt.fname,
  638. "fileType": tt.ft,
  639. "rollingInterval": 0,
  640. "rollingCount": 1,
  641. "rollingNamePattern": "none",
  642. "hasHeader": true,
  643. "format": "delimited",
  644. "compression": tt.compress,
  645. })
  646. if err != nil {
  647. t.Fatal(err)
  648. }
  649. mockclock.ResetClock(10)
  650. err = sink.Open(ctx)
  651. if err != nil {
  652. t.Fatal(err)
  653. }
  654. c := mockclock.GetMockClock()
  655. for i := 0; i < 3; i++ {
  656. c.Add(450 * time.Millisecond)
  657. m := map[string]interface{}{"key": "value" + strconv.Itoa(i), "ts": c.Now().UnixMilli()}
  658. if err := sink.Collect(vCtx, m); err != nil {
  659. t.Errorf("unexpected error: %s", err)
  660. }
  661. }
  662. c.After(2000 * time.Millisecond)
  663. if err = sink.Close(ctx); err != nil {
  664. t.Errorf("unexpected close error: %s", err)
  665. }
  666. // Should write to 2 files
  667. for i := 0; i < 3; i++ {
  668. // Read the contents of the temporary file and check if they match the collected items
  669. var fn string
  670. if tt.compress != "" {
  671. fn = fmt.Sprintf("test_%s_%s_%d.dd", tt.ft, tt.compress, 460+450*i)
  672. } else {
  673. fn = fmt.Sprintf("test_%s_%d.dd", tt.ft, 460+450*i)
  674. }
  675. contents, err := os.ReadFile(fn)
  676. if err != nil {
  677. t.Fatal(err)
  678. }
  679. if tt.compress != "" {
  680. decompressor, _ := compressor.GetDecompressor(tt.compress)
  681. contents, err = decompressor.Decompress(contents)
  682. if err != nil {
  683. t.Errorf("%v", err)
  684. }
  685. }
  686. if !reflect.DeepEqual(contents, tt.contents[i]) {
  687. t.Errorf("\nexpected\t %q \nbut got\t\t %q", tt.contents[i], string(contents))
  688. }
  689. }
  690. })
  691. }
  692. }
  693. func TestFileSinkReopen(t *testing.T) {
  694. // Remove existing files
  695. err := filepath.Walk(".", func(path string, info os.FileInfo, err error) error {
  696. if err != nil {
  697. return err
  698. }
  699. if filepath.Ext(path) == ".log" {
  700. fmt.Println("Deleting file:", path)
  701. return os.Remove(path)
  702. }
  703. return nil
  704. })
  705. if err != nil {
  706. t.Fatal(err)
  707. }
  708. conf.IsTesting = true
  709. tmpfile, err := os.CreateTemp("", "reopen.log")
  710. if err != nil {
  711. t.Fatal(err)
  712. }
  713. defer os.Remove(tmpfile.Name())
  714. // Create a stream context for testing
  715. contextLogger := conf.Log.WithField("rule", "testRollingCount")
  716. ctx := context.WithValue(context.Background(), context.LoggerKey, contextLogger)
  717. tf, _ := transform.GenTransform("", "json", "", "", "", []string{})
  718. vCtx := context.WithValue(ctx, context.TransKey, tf)
  719. sink := &fileSink{}
  720. err = sink.Configure(map[string]interface{}{
  721. "path": tmpfile.Name(),
  722. "fileType": LINES_TYPE,
  723. "format": "json",
  724. "rollingNamePattern": "none",
  725. })
  726. if err != nil {
  727. t.Fatal(err)
  728. }
  729. err = sink.Open(vCtx)
  730. if err != nil {
  731. t.Fatal(err)
  732. }
  733. // Test collecting a map item
  734. m := map[string]interface{}{"key": "value1"}
  735. if err := sink.Collect(vCtx, m); err != nil {
  736. t.Errorf("unexpected error: %s", err)
  737. }
  738. sink.Close(vCtx)
  739. exp := []byte(`{"key":"value1"}`)
  740. contents, err := os.ReadFile(tmpfile.Name())
  741. if err != nil {
  742. t.Fatal(err)
  743. }
  744. if !reflect.DeepEqual(contents, exp) {
  745. t.Errorf("\nexpected\t %q \nbut got\t\t %q", string(exp), string(contents))
  746. }
  747. sink = &fileSink{}
  748. err = sink.Configure(map[string]interface{}{
  749. "path": tmpfile.Name(),
  750. "fileType": LINES_TYPE,
  751. "hasHeader": true,
  752. "format": "json",
  753. "rollingNamePattern": "none",
  754. })
  755. if err != nil {
  756. t.Fatal(err)
  757. }
  758. err = sink.Open(vCtx)
  759. if err != nil {
  760. t.Fatal(err)
  761. }
  762. // Test collecting another map item
  763. m = map[string]interface{}{"key": "value2"}
  764. if err := sink.Collect(vCtx, m); err != nil {
  765. t.Errorf("unexpected error: %s", err)
  766. }
  767. if err = sink.Close(vCtx); err != nil {
  768. t.Errorf("unexpected close error: %s", err)
  769. }
  770. exp = []byte(`{"key":"value2"}`)
  771. contents, err = os.ReadFile(tmpfile.Name())
  772. if err != nil {
  773. t.Fatal(err)
  774. }
  775. if !reflect.DeepEqual(contents, exp) {
  776. t.Errorf("\nexpected\t %q \nbut got\t\t %q", string(exp), string(contents))
  777. }
  778. }