file.go 6.7 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336
  1. package rx
  2. import (
  3. "os"
  4. "io"
  5. "fmt"
  6. "time"
  7. "bufio"
  8. "errors"
  9. "math/big"
  10. "io/ioutil"
  11. "path/filepath"
  12. "kumachan/standalone/util"
  13. )
  14. type File struct {
  15. raw *os.File
  16. worker *Worker
  17. }
  18. type FileState struct {
  19. Name string
  20. Size *big.Int
  21. Mode uint32
  22. IsDir bool
  23. ModTime time.Time
  24. }
  25. func FileStateFromInfo(info os.FileInfo) FileState {
  26. return FileState {
  27. Name: info.Name(),
  28. Size: big.NewInt(info.Size()),
  29. Mode: uint32(info.Mode()),
  30. IsDir: info.IsDir(),
  31. ModTime: info.ModTime(),
  32. }
  33. }
  34. func FileFrom(raw *os.File) File {
  35. return File {
  36. raw: raw,
  37. worker: CreateWorker(),
  38. }
  39. }
  40. func OpenReadOnly(path string) Observable {
  41. return Open(path, os.O_RDONLY, 0)
  42. }
  43. func OpenReadWrite(path string) Observable {
  44. return Open(path, os.O_RDWR, 0)
  45. }
  46. func OpenReadWriteCreate(path string, perm os.FileMode) Observable {
  47. return Open(path, os.O_RDWR | os.O_CREATE, perm)
  48. }
  49. func OpenOverwrite(path string, perm os.FileMode) Observable {
  50. return Open(path, os.O_WRONLY | os.O_APPEND | os.O_CREATE | os.O_TRUNC, perm)
  51. }
  52. func OpenAppend(path string, perm os.FileMode) Observable {
  53. return Open(path, os.O_WRONLY | os.O_APPEND | os.O_CREATE, perm)
  54. }
  55. func Open(path string, flag int, perm os.FileMode) Observable {
  56. return NewGoroutine(func(sender Sender) {
  57. if sender.Context().AlreadyCancelled() {
  58. return
  59. }
  60. raw, err := os.OpenFile(path, flag, perm)
  61. if err != nil {
  62. sender.Error(err)
  63. return
  64. }
  65. var f = File {
  66. raw: raw,
  67. worker: CreateWorker(),
  68. }
  69. sender.Next(f)
  70. sender.Complete()
  71. sender.Context().WaitDispose(func() {
  72. _ = raw.Close()
  73. f.worker.Dispose()
  74. })
  75. })
  76. }
  77. func (f File) Close() Observable {
  78. return NewQueued(f.worker, func() (Object, bool) {
  79. _ = f.raw.Close()
  80. f.worker.Dispose()
  81. return nil, true
  82. })
  83. }
  84. func (f File) State() Observable {
  85. return NewQueued(f.worker, func() (Object, bool) {
  86. var info, err = f.raw.Stat()
  87. if err != nil {
  88. return err, false
  89. } else {
  90. return FileStateFromInfo(info), true
  91. }
  92. })
  93. }
  94. func (f File) Read(amount uint) Observable {
  95. return NewQueued(f.worker, func() (Object, bool) {
  96. var buf = make([] byte, amount)
  97. var n, err = f.raw.Read(buf)
  98. if err != nil {
  99. return err, false
  100. } else {
  101. var result = buf[:n]
  102. return result, true
  103. }
  104. })
  105. }
  106. func (f File) Write(data ([] byte)) Observable {
  107. return NewQueued(f.worker, func() (Object, bool) {
  108. var _, err = f.raw.Write(data)
  109. if err != nil {
  110. return err, false
  111. } else {
  112. return nil, true
  113. }
  114. })
  115. }
  116. func (f File) SeekStart(offset int64) Observable {
  117. return NewQueued(f.worker, func() (Object, bool) {
  118. var new_offset, err = f.raw.Seek(offset, io.SeekStart)
  119. if err != nil {
  120. return err, false
  121. }
  122. return new_offset, true
  123. })
  124. }
  125. func (f File) SeekDelta(delta int64) Observable {
  126. return NewQueued(f.worker, func() (Object, bool) {
  127. var new_offset, err = f.raw.Seek(delta, io.SeekCurrent)
  128. if err != nil {
  129. return err, false
  130. }
  131. return new_offset, true
  132. })
  133. }
  134. func (f File) SeekEnd(offset int64) Observable {
  135. return NewQueued(f.worker, func() (Object, bool) {
  136. var new_offset, err = f.raw.Seek(offset, io.SeekEnd)
  137. if err != nil {
  138. return err, false
  139. }
  140. return new_offset, true
  141. })
  142. }
  143. func (f File) ReadChar() Observable {
  144. return NewQueued(f.worker, func() (Object, bool) {
  145. var char rune
  146. var _, err = fmt.Fscanf(f.raw, "%c", &char)
  147. if err != nil {
  148. return err, false
  149. }
  150. return char, true
  151. })
  152. }
  153. func (f File) WriteChar(char rune) Observable {
  154. return NewQueued(f.worker, func() (Object, bool) {
  155. var _, err = fmt.Fprintf(f.raw, "%c", char)
  156. if err != nil {
  157. return err, false
  158. }
  159. return nil, true
  160. })
  161. }
  162. func (f File) ReadRunes() Observable {
  163. return NewQueued(f.worker, func() (Object, bool) {
  164. var buf = make([] rune, 0)
  165. for {
  166. var char rune
  167. var _, err = fmt.Fscanf(f.raw, "%c", &char)
  168. if err != nil { return err, false }
  169. if char != ' ' && char != '\n' {
  170. buf = append(buf, char)
  171. } else {
  172. return buf, true
  173. }
  174. }
  175. })
  176. }
  177. func (f File) ReadString() Observable {
  178. return f.ReadRunes().Map(func(runes Object) Object {
  179. return string(runes.([] rune))
  180. })
  181. }
  182. func (f File) ReadLineRunes() Observable {
  183. return NewQueued(f.worker, func() (Object, bool) {
  184. var str, err = util.WellBehavedReadLine(f.raw)
  185. if err != nil {
  186. return err, false
  187. }
  188. return str, true
  189. })
  190. }
  191. func (f File) ReadLine() Observable {
  192. return f.ReadLineRunes().Map(func(runes Object) Object {
  193. return string(runes.([] rune))
  194. })
  195. }
  196. func (f File) ReadAll() Observable {
  197. return NewQueued(f.worker, func() (Object, bool) {
  198. var bytes, err = ioutil.ReadAll(f.raw)
  199. if err != nil {
  200. return err, false
  201. }
  202. return bytes, true
  203. })
  204. }
  205. func (f File) WriteString(str string) Observable {
  206. return NewQueued(f.worker, func() (Object, bool) {
  207. var _, err = fmt.Fprint(f.raw, str)
  208. if err != nil {
  209. return err, false
  210. }
  211. return nil, true
  212. })
  213. }
  214. func (f File) WriteLine(str string) Observable {
  215. return NewQueued(f.worker, func() (Object, bool) {
  216. var _, err = fmt.Fprintln(f.raw, str)
  217. if err != nil {
  218. return err, false
  219. }
  220. return nil, true
  221. })
  222. }
  223. func (f File) ReadLinesRuneSlices() Observable {
  224. // emits rune slices
  225. return NewGoroutine(func(s Sender) {
  226. f.worker.Do(func() {
  227. var buffered = bufio.NewReader(f.raw)
  228. for {
  229. if s.Context().AlreadyCancelled() {
  230. return
  231. }
  232. var line, err = util.WellBehavedReadLine(buffered)
  233. if err != nil {
  234. if err == io.EOF {
  235. s.Complete()
  236. return
  237. } else {
  238. s.Error(err)
  239. return
  240. }
  241. }
  242. s.Next(line)
  243. }
  244. })
  245. })
  246. }
  247. func (f File) ReadLines() Observable {
  248. return f.ReadLinesRuneSlices().Map(func(runes Object) Object {
  249. return string(runes.([] rune))
  250. })
  251. }
  252. type FileItem struct {
  253. Path string
  254. State FileState
  255. }
  256. func WalkDir(root string) Observable {
  257. return NewGoroutine(func(sender Sender) {
  258. var err = filepath.Walk(root, func(path string, info os.FileInfo, err error) error {
  259. if sender.Context().AlreadyCancelled() {
  260. return errors.New("operation cancelled")
  261. }
  262. sender.Next(FileItem {
  263. Path: path,
  264. State: FileStateFromInfo(info),
  265. })
  266. return err
  267. })
  268. if err != nil {
  269. sender.Error(err)
  270. return
  271. }
  272. sender.Complete()
  273. })
  274. }
  275. func ListDir(dir_path string) Observable {
  276. return NewGoroutine(func(sender Sender) {
  277. var err = filepath.Walk(dir_path, func(path string, info os.FileInfo, err error) error {
  278. if sender.Context().AlreadyCancelled() {
  279. return errors.New("operation cancelled")
  280. }
  281. if path != dir_path {
  282. sender.Next(FileItem {
  283. Path: path,
  284. State: FileStateFromInfo(info),
  285. })
  286. if info.IsDir() && path != dir_path {
  287. return filepath.SkipDir
  288. } else {
  289. return nil
  290. }
  291. } else {
  292. return nil
  293. }
  294. })
  295. if err != nil {
  296. sender.Error(err)
  297. return
  298. }
  299. sender.Complete()
  300. })
  301. }